smaple.tr
sharding

Database Sharding and Partitioning: Horizontal Scaling Strategies [2026]

Mehmet Kurtipek
April 5, 2026
10 min read
sharding
partitioning
database scaling
horizontal scaling
shard key
consistent hashing

Every database hits a single-server ceiling. As data volume grows past several hundred gigabytes and query concurrency increases past what one server can handle, two scaling strategies emerge: vertical scaling (upgrading to a larger server) and horizontal scaling (distributing data across multiple servers). Partitioning and sharding are the techniques that implement horizontal scaling — and both are frequently misapplied, adding operational complexity without the performance gains that motivated the change.

This guide covers database sharding and partitioning from the decision criteria through production implementation: the difference between partitioning and sharding, the four partitioning strategies and when each applies, shard key selection (the decision that most often goes wrong), consistent hashing for even distribution, PostgreSQL declarative partitioning, MongoDB sharding, and the operational patterns for managing distributed data at scale.

Database Sharding vs Partitioning: The Core Distinction

Partitioning divides a table into multiple physical segments within a single database server. The database engine handles routing transparently — queries look the same whether the table is partitioned or not. Partitioning improves performance through partition pruning (the query planner reads only relevant partitions) and improves manageability (old partitions can be dropped instantly with DROP TABLE).

Sharding distributes data across multiple independent database servers (shards). Each shard owns a subset of the data. The application layer (or a routing proxy) must direct queries to the correct shard. Sharding enables true horizontal scaling — read and write capacity scales linearly with the number of shards.

Start with partitioning. Most applications never need sharding. Partitioning addresses the most common scaling pain points: query performance on large tables, fast data lifecycle management (deleting old data), and parallel query execution. Sharding introduces significant operational complexity — multi-shard transactions, cross-shard queries, shard rebalancing — that partitioning avoids entirely.

When sharding becomes necessary:

  • Write throughput exceeds a single server's capacity
  • Dataset size exceeds the maximum vertical upgrade available (e.g., 10TB+ for cloud instances)
  • Geographic data distribution requirements demand data residency in specific regions
  • Read replicas cannot keep up with read traffic and adding more replicas creates unacceptable replication lag

Partitioning Strategies

PostgreSQL supports four declarative partitioning strategies. The right strategy depends on the primary access pattern for the data.

Range Partitioning

Best for time-series operational data where queries naturally filter by date ranges.

CREATE TABLE events (
    id BIGSERIAL,
    entity_id INTEGER NOT NULL,
    event_type VARCHAR(50),
    payload JSONB,
    created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
) PARTITION BY RANGE (created_at);

-- Monthly partitions
CREATE TABLE events_2026_01 PARTITION OF events
    FOR VALUES FROM ('2026-01-01') TO ('2026-02-01');

CREATE TABLE events_2026_02 PARTITION OF events
    FOR VALUES FROM ('2026-02-01') TO ('2026-03-01');

CREATE TABLE events_2026_03 PARTITION OF events
    FOR VALUES FROM ('2026-03-01') TO ('2026-04-01');

Partition pruning in action:

-- This query only reads events_2026_03 — the planner prunes all other partitions
SELECT * FROM events
WHERE created_at >= '2026-03-01' AND created_at < '2026-04-01'
  AND entity_id = 42;

Partition maintenance:

-- Drop old partition (instant — no DELETE, no VACUUM needed)
DROP TABLE events_2025_01;

-- Attach a pre-populated partition (for data migrations)
ALTER TABLE events ATTACH PARTITION events_2026_04
    FOR VALUES FROM ('2026-04-01') TO ('2026-05-01');

Hash Partitioning

Even distribution of data across a fixed number of partitions. Best for tables where no natural range key exists and even load distribution is the priority.

CREATE TABLE users (
    id BIGSERIAL,
    email VARCHAR(255) UNIQUE NOT NULL,
    created_at TIMESTAMPTZ DEFAULT NOW()
) PARTITION BY HASH (id);

-- 4 partitions with modulus=4
CREATE TABLE users_p0 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 0);
CREATE TABLE users_p1 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 1);
CREATE TABLE users_p2 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 2);
CREATE TABLE users_p3 PARTITION OF users FOR VALUES WITH (modulus 4, remainder 3);

Hash partitioning limitation: Adding partitions requires redistributing all data — there is no cheap way to go from 4 to 5 partitions. Plan the partition count for expected 3–5 year growth.

List Partitioning

Partition by discrete values. Best for data that naturally groups by a finite set of categories (country, region, tenant tier).

CREATE TABLE orders (
    id BIGSERIAL,
    region VARCHAR(20) NOT NULL,
    customer_id INTEGER,
    total_amount DECIMAL(12, 2),
    created_at TIMESTAMPTZ DEFAULT NOW()
) PARTITION BY LIST (region);

CREATE TABLE orders_north PARTITION OF orders FOR VALUES IN ('US', 'CA', 'MX');
CREATE TABLE orders_europe PARTITION OF orders FOR VALUES IN ('DE', 'FR', 'GB', 'NL');
CREATE TABLE orders_apac PARTITION OF orders FOR VALUES IN ('AU', 'JP', 'SG', 'IN');
CREATE TABLE orders_other PARTITION OF orders DEFAULT;  -- Catch-all partition

Composite Partitioning (Sub-Partitioning)

Partition by one strategy at the top level, then sub-partition by another strategy.

-- Primary: range by year; secondary: list by region
CREATE TABLE sales (
    id BIGSERIAL,
    sale_date DATE NOT NULL,
    region VARCHAR(20) NOT NULL,
    amount DECIMAL(12, 2)
) PARTITION BY RANGE (sale_date);

CREATE TABLE sales_2026 PARTITION OF sales
    FOR VALUES FROM ('2026-01-01') TO ('2027-01-01')
    PARTITION BY LIST (region);

CREATE TABLE sales_2026_north PARTITION OF sales_2026
    FOR VALUES IN ('US', 'CA');
CREATE TABLE sales_2026_europe PARTITION OF sales_2026
    FOR VALUES IN ('DE', 'FR', 'GB');

Sharding Architecture

Shard Key Selection

The shard key is the column (or combination of columns) that determines which shard a given row lives on. A poor shard key choice produces hot shards (all traffic goes to one shard), cross-shard queries (most queries must hit multiple shards), or requires a full cluster rebuild to change.

Good shard key properties:

High cardinality: Many distinct values. A boolean field or a status column with 5 values is a terrible shard key — you cannot distribute 1TB of data across more shards than there are distinct values.

Uniform write distribution: Avoid monotonically increasing keys (like auto-increment IDs or timestamps) as the primary shard key. All new writes go to the latest shard, creating a hotspot. Hash the key or use a different dimension.

Query locality: Ideally, most queries include the shard key in their WHERE clause. A query without the shard key must scatter-gather across all shards, eliminating the performance benefit.

Multi-tenant SaaS pattern: tenant_id as the shard key is a common production choice. Each tenant's data lives on one shard, so all tenant-scoped queries go to a single shard. Cross-tenant queries are rare in SaaS. The risk is large-tenant hotspots — one huge tenant can overload a shard.

Consistent Hashing

Consistent hashing is the algorithm that maps shard keys to shards and handles shard addition/removal with minimal data movement.

The concept: imagine a ring with 0–360 degrees. Each shard is assigned a position on the ring. A key is hashed to a position on the ring, then assigned to the first shard clockwise from that position.

Advantage over simple modulo hashing (key % N): When you add a shard (change N from 4 to 5), modulo hashing remaps most keys. Consistent hashing remaps only approximately 1/N of keys when adding a shard.

import hashlib
import bisect

class ConsistentHashRing:
    def __init__(self, nodes, virtual_nodes=150):
        self.ring = {}
        self.sorted_keys = []

        for node in nodes:
            for i in range(virtual_nodes):
                key = self._hash(f"{node}:{i}")
                self.ring[key] = node
                bisect.insort(self.sorted_keys, key)

    def _hash(self, key: str) -> int:
        return int(hashlib.md5(key.encode()).hexdigest(), 16)

    def get_node(self, key: str) -> str:
        if not self.ring:
            return None
        hash_key = self._hash(key)
        idx = bisect.bisect(self.sorted_keys, hash_key) % len(self.sorted_keys)
        return self.ring[self.sorted_keys[idx]]

# Usage
ring = ConsistentHashRing(['shard-1', 'shard-2', 'shard-3', 'shard-4'])
shard = ring.get_node('tenant_id:12345')  # Returns consistent shard assignment

Virtual nodes: Each physical shard is represented by multiple virtual nodes on the ring (150 in the example above). This ensures even distribution even with a small number of physical shards.

MongoDB Sharding

MongoDB's sharding implementation is the most mature NoSQL sharding solution.

// Enable sharding on a database
sh.enableSharding("myapp")

// Shard a collection on tenant_id (hash-based for even distribution)
sh.shardCollection("myapp.orders", { tenant_id: "hashed" })

// Or range-based (better for range queries, hotspot risk)
sh.shardCollection("myapp.events", { tenant_id: 1, created_at: 1 })

// Check sharding status
sh.status()

MongoDB shard key requirements:

  • The shard key must be present in every document
  • The shard key must be part of every index
  • Shard key cannot be changed after the collection is sharded (pre-7.0)

Zone sharding for data residency:

// Assign specific key ranges to specific shards (e.g., for data sovereignty)
sh.addShardToZone("shard-eu-west-1", "EU")
sh.addShardToZone("shard-us-east-1", "US")

sh.updateZoneKeyRange("myapp.users",
  { region: "DE" },
  { region: "DE\uffff" },
  "EU"
)

Cross-Shard Query Patterns

Cross-shard queries — queries that must read from multiple shards — are the main performance concern in sharded architectures. Three patterns minimize their frequency.

Denormalization for co-location: Store copies of related data on the same shard. A user's complete profile, their orders, and their activity history should all be on the same shard (same user_id shard key), even if it means duplicating some data.

Scatter-gather with aggregation: When a cross-shard query is unavoidable, send it to all shards in parallel and aggregate results in the application layer. MongoDB's mongos router handles this automatically for sharded aggregations.

Secondary indexes in application layer: For non-shard-key lookups, maintain a lookup table (in a separate database) that maps the lookup key to the shard key. email → user_id lookup enables routing a login query (which knows email, not user_id) to the correct shard without scanning all shards.

Partitioning Best Practices

Partition key selection for range partitions: The most common mistake in range partitioning is choosing a partition key that creates hot partitions. If all recent activity goes to the current month's partition, that partition becomes a write bottleneck while historical partitions receive no load. For write-heavy current data, combine range partitioning (by month) with hash sub-partitioning (by entity ID) to distribute writes across multiple physical segments within each time range.

Partition pruning verification: Always verify the query planner uses partition pruning for your production queries:

EXPLAIN (ANALYZE, PARTITIONS)
SELECT * FROM events
WHERE created_at >= '2026-03-01' AND created_at < '2026-04-01';
-- Should show: Partitions selected: 1 out of 24
-- If it shows "Partitions selected: 24 out of 24", pruning is not working

Partition pruning requires that the partition key column appears in the WHERE clause with a literal or parameterized value. If the partition key is cast or wrapped in a function call, pruning is suppressed.

Partition automation: Creating partitions manually before data arrives is fragile — a missing partition causes insert failures. PostgreSQL 16 does not yet support automatic partition creation, but pg_partman is a widely-used extension that automates time-based partition maintenance (creating future partitions and dropping expired ones on schedule).

Operational Considerations

Shard rebalancing: When data grows unevenly, shards can be rebalanced. MongoDB's balancer moves chunks between shards automatically. PostgreSQL partitioning does not auto-rebalance — this is a manual operation.

Backup complexity: Backing up a sharded cluster requires coordinated snapshots across all shards. Partial backups that are not time-consistent across shards produce corrupt restore states. Use the database vendor's recommended backup procedures.

Schema changes: In sharded systems, schema migrations must be applied to each shard. Phased migrations (backward-compatible schema changes deployed before application code changes) are required to avoid downtime during schema rollouts.

Conclusion

Database sharding and partitioning solve different problems at different scales. Partitioning — range by date, hash by ID, list by category — addresses the most common performance pain points for large tables without introducing distributed system complexity. Sharding is the appropriate solution when write throughput or data volume genuinely exceeds single-server capacity.

The most common sharding mistake is premature adoption. Shard key selection is permanent in most systems, and a bad shard key that creates hotspots or requires scatter-gather for common queries eliminates the performance benefit and cannot be easily corrected without rebuilding the cluster.


Author: Smart Maple Database Engineering Team Updated: April 2026

Related Articles

August 11, 2026

MLOps Guide: Taking Machine Learning Models to Production [2026]

87% of machine learning models built by data science teams never reach production. The models work — they pass cross-validation, they score well on holdout sets, they demonstrate genuine predictive value. The problem is not the modeling. The problem is everything that happens between a notebook experiment and a reliable, monitored, production system. MLOps is the discipline that closes that gap. This guide covers the full MLOps stack: maturity levels, tooling choices (MLflow, DVC, Kubeflow

Read More
August 10, 2026

LLM Fine-Tuning Guide: Custom Model Training with LoRA and QLoRA [2026]

General-purpose LLMs are impressive. They can write code, summarize documents, answer questions, and translate between languages with reasonable accuracy. But "reasonable" is not good enough when your application requires consistent output format, domain-specific terminology, a particular tone, or behavior that the base model was never trained to exhibit. That gap is where fine-tuning matters. Fine-tuning updates a model's weights on your specific data, changing how the model behaves — not

Read More
August 9, 2026

Computer Vision Applications: Object Detection, OCR, and Industrial AI [2026]

Computer vision has moved well past the research phase. The models are trained, the frameworks are mature, the hardware is accessible, and the use cases are generating measurable returns. What was a specialized capability requiring deep expertise in 2018 is now deployable infrastructure — if you know which component to reach for and where the real complexity lives. This guide covers computer vision applications across industrial, medical, logistics, and document processing domains. It expl

Read More