Data partitioning divides a dataset into smaller parts so storage, computation, retention or access can be managed more efficiently. Sharding is a form of horizontal partitioning in which different subsets of records are distributed across independent nodes or database instances. The central design question is not merely how to split the data, but which split preserves locality, avoids hotspots, supports growth and still allows the organisation to find and recombine what it needs.
Partitioning turns one large data problem into several smaller ones. A poor partition key turns those smaller problems into a coordination problem.
As datasets grow, one machine or one storage unit can become a bottleneck. Partitioning can spread data across files, disks, nodes, regions or logical boundaries. The gain is scale and locality. The cost is that some questions now cross boundaries and require routing, coordination and rebalancing.
ARTICLE ID: DATA.MANAGEMENT.052
Canonical function: physical and logical distribution of large datasets across manageable boundaries
Owner boundary: this article owns partitioning and sharding. Database Management and Transactional Integrity owns transactional correctness; Data Lakes and Lakehouses owns analytical object-storage architecture; Data Architecture owns overall system structure.
The Simple Answer
A scalable partitioning route is:
Understand Access Pattern → Choose Partition Boundary → Choose Distribution Key → Route Writes → Route Reads → Observe Skew → Rebalance → Reconcile → Evolve Without Losing Identity
The partition key is one of the most consequential choices because it determines which records live together and which operations require cross-partition work.
Horizontal Partitioning
Horizontal partitioning divides rows across partitions while keeping broadly the same columns.
Examples:
- customers A–F in one partition and G–L in another;
- orders split by customer ID hash;
- events split by month;
- tenants split by organisation ID.
Vertical Partitioning
Vertical partitioning divides columns or functional groups of fields.
A customer profile can keep common operational fields in one store while large documents, sensitive attributes or rarely used analytical fields live elsewhere.
Vertical partitioning changes access paths and should preserve stable entity identity across the split.
Sharding
Sharding distributes horizontal partitions across separate database or service nodes so capacity can scale beyond one node.
Each shard usually owns a subset of keys. A routing layer decides which shard should receive a request.
Partition Key
The partition key determines how records are assigned.
- customer ID;
- tenant ID;
- date;
- region;
- product ID;
- hash of a stable identifier.
A good key spreads load while keeping related operations local enough for efficient reads and writes.
Choose for Access Patterns, Not Aesthetics
A key that looks evenly distributed may be poor if most queries need data from several shards. A key with perfect locality may be poor if one entity becomes far larger than all others.
The design should begin with actual read, write and growth patterns.
Range Partitioning
Range partitioning places contiguous key ranges together.
- dates by month;
- customer IDs 1–1,000,000;
- alphabetic surname ranges;
- geographic coordinate ranges.
Range partitions support efficient range scans but can create hotspots when new writes concentrate in one current range.
Time Partitioning
Time is a common range key for event and analytical data.
Benefits include:
- partition pruning for time-bounded queries;
- easy lifecycle management;
- efficient dropping or archiving of old periods;
- bounded backfills.
The current time partition can become a write hotspot under heavy ingestion.
Hash Partitioning
Hash partitioning applies a deterministic hash to a key and maps the result to a partition.
It often distributes records more evenly than natural ranges, but destroys locality for range scans on the original key.
List Partitioning
List partitioning assigns explicit categories to partitions.
Examples include country, business unit, product family or data classification.
List partitions are easy to understand but can become unbalanced when categories grow unevenly.
Composite Partitioning
Composite strategies combine keys or methods, such as partitioning by month and then hashing by customer.
This can preserve useful time pruning while spreading current-period load.
Locality
Locality means keeping data that is commonly accessed together in the same partition or nearby storage.
Good locality reduces network transfer and cross-shard coordination.
Entity Locality
Partitioning by customer or tenant can keep most entity-specific operations on one shard.
This is powerful for multi-tenant systems, but very large tenants can become hotspots.
Temporal Locality
Recent data is often queried more frequently than old data. Time partitioning can keep hot recent partitions on fast storage and older partitions on cheaper tiers.
Geographic Locality
Some systems place data near users or operations by region to reduce latency or satisfy residency constraints.
Geographic partitioning can complicate global queries and user movement between regions.
See Data Sovereignty, Residency and Jurisdiction.
Hotspots
A hotspot occurs when one partition receives substantially more traffic, storage growth or computation than others.
Causes include:
- one extremely large tenant;
- sequential keys concentrating writes;
- popular products;
- current-time partitions;
- uneven geographic activity;
- poor hash choice;
- celebrity or viral entities.
Skew
Partition skew means data or workload is distributed unevenly.
Skew can be measured by storage size, rows, write rate, read rate, CPU, memory or network traffic. A partition balanced by rows can still be unbalanced by workload.
Cardinality
Partition keys usually need enough distinct values to distribute load effectively.
A boolean key creates at most two partitions and is rarely useful for large sharding schemes.
High Cardinality Is Not Enough
A customer ID has high cardinality, but one enterprise customer can still create a hotspot if most workload belongs to that customer.
Evaluate frequency distribution as well as distinct count.
Routing
A sharded system needs a route from logical key to physical shard.
- client-side routing;
- proxy or gateway routing;
- directory lookup;
- hash calculation;
- metadata service;
- database-native routing.
Routing metadata becomes critical infrastructure and should be versioned and highly available.
Directory-Based Sharding
A directory maps entities or key ranges to shards explicitly.
This supports flexible movement during rebalancing but creates dependence on the routing directory.
Hash-Based Routing
Hash routing calculates the destination from a stable key. It avoids a large central mapping table but can make shard-count changes expensive unless the scheme supports controlled remapping.
Consistent Hashing
Consistent-hashing techniques are designed so adding or removing nodes remaps only part of the keyspace rather than every key.
They can reduce movement during rebalancing, though real systems often add virtual nodes, weights or directories to improve load balance.
Rebalancing
Rebalancing moves partitions or key ranges when capacity, workload or topology changes.
A safe rebalancing process must preserve availability and avoid losing writes during movement.
Online Rebalancing
Online rebalancing moves data while the system continues serving traffic.
Typical concerns include:
- copying existing state;
- capturing concurrent writes;
- switching routing;
- reconciling source and destination;
- retiring old copies;
- handling rollback or ambiguous outcomes.
See Data Synchronisation and Reconciliation.
Shard Splitting
A large shard can be split into smaller ranges or hash buckets when it reaches capacity.
Designing for split-friendly keys reduces emergency migration work later.
Shard Merging
Underused shards may be merged to reduce operational cost. Merging should preserve stable identities and routing correctness just like splitting.
Replication Is Not Partitioning
Partitioning divides data. Replication creates additional copies.
A system can use both: each shard owns part of the keyspace, and each shard is replicated for availability.
Primary and Replica
One common design assigns a primary replica for writes and additional replicas for resilience or read scale.
Replication lag can make reads stale, so the routing layer may need consistency-aware choices.
Cross-Shard Queries
A query that touches many shards must fan out, gather results and combine them.
This increases latency, network traffic and coordination cost. Global analytical workloads are often better served by dedicated analytical copies than by repeatedly scanning operational shards.
Scatter-Gather
Scatter-gather sends a request to several shards and merges responses.
The slowest shard can dominate the result. A single failing shard can create partial or failed queries depending on the contract.
Cross-Shard Joins
Joins are cheapest when related records are colocated. Cross-shard joins may require moving large datasets or performing distributed query coordination.
Partitioning should therefore consider common join keys as well as write distribution.
Cross-Shard Transactions
Transactions confined to one shard are usually simpler. Multi-shard transactions require coordination and can reduce availability or increase latency.
Domain design can sometimes keep high-integrity operations local by partitioning related entities together.
Global Uniqueness
Unique constraints become harder across independent shards.
Options include:
- globally generated IDs;
- partition-prefixed IDs;
- central allocation services;
- probabilistically unique identifiers;
- eventual duplicate detection with repair.
The choice depends on whether uniqueness must be guaranteed synchronously.
Global Secondary Indexes
If data is partitioned by customer ID but users search by email, a global secondary index may map email to the owning shard.
That index becomes another distributed data product requiring synchronisation, uniqueness rules and recovery.
Denormalisation Across Shards
Systems sometimes duplicate reference or lookup data onto each shard to avoid cross-shard joins.
This improves locality but creates synchronisation obligations. Derived copies should preserve source authority and version.
Multi-Tenancy
Tenant ID is a common partition key because it provides organisational locality and can simplify access control.
However, tenants often differ dramatically in size. Large tenants may require sub-sharding or dedicated placement.
Noisy Neighbours
One tenant or workload can consume disproportionate CPU, memory or I/O and degrade others sharing the shard.
Workload isolation, quotas and tenant-aware routing can reduce noisy-neighbour effects.
Data Sovereignty and Geo-Sharding
Geo-sharding places records in designated regions according to residency, latency or operational requirements.
Routing must account for user movement, cross-border queries, backups and metadata. A regional shard map itself can become sensitive operational information.
Partition Pruning
Partition pruning lets query engines skip partitions that cannot contain relevant records.
If a table is partitioned by month and a query requests one month, the engine can avoid scanning other periods when filters align with the partition key.
Too Many Partitions
Very small partitions create metadata, planning and file-management overhead.
A partition per minute may be excessive when a daily partition would provide the same pruning value with far lower management cost.
Too Few Partitions
Very large partitions reduce pruning and make maintenance, rebalancing or recovery more expensive.
Partition size should reflect both query behaviour and operational manageability.
Small Files in Analytical Partitions
Analytical systems can suffer when one logical partition contains thousands of tiny files. Planning and metadata overhead may dominate actual data scanning.
Compaction can combine small files while preserving partition semantics.
See Data Lakes and Lakehouses.
Partition Evolution
Access patterns change. A key that worked at one scale may become unsuitable later.
Partition strategy should therefore have an evolution path rather than being treated as permanent physical truth.
Changing the Partition Key
Changing a partition key often requires redistributing most or all records.
A safe migration can use:
- new target shards;
- bulk copy;
- change capture;
- dual reads or controlled dual writes;
- reconciliation;
- routing cutover;
- old-shard retirement.
See Data Migration and Legacy Modernisation.
Consistent Identity Through Repartitioning
Physical location can change without changing logical identity.
Consumers should not use shard location as the permanent business identifier for an entity.
Failure Domains
Partitioning can isolate failure. One shard may fail while others remain healthy.
This is beneficial only if routing, replication and application behaviour can tolerate partial availability.
Partial Availability
A sharded system may serve most users while one shard is unavailable.
Global summaries should not silently treat the missing shard as zero. Partial-state indicators belong in receiver-facing contracts.
Backups
Each shard requires backup and recovery coverage. A global recovery plan must know whether shards can be restored independently or require a coordinated consistency point.
See Data Backup, Recovery and Resilience.
Reconciliation After Rebalancing
Moving a partition should end with independent evidence that source and destination contain the expected keyspace, counts, control totals or checksums.
“Copy job succeeded” is not sufficient proof.
Observability
Useful signals include:
- rows and bytes per partition;
- read and write rate;
- latency by shard;
- CPU and memory;
- cross-shard query rate;
- hot-key concentration;
- rebalancing backlog;
- replication lag;
- routing failures;
- partial-availability incidents.
Skew Metrics
Compare largest, smallest, median and percentile shard loads rather than relying on average utilisation.
A healthy average can hide one overloaded partition near failure.
Testing
Tests should cover:
- routing correctness;
- boundary keys;
- hotspot scenarios;
- node loss;
- replica lag;
- cross-shard queries;
- rebalancing under concurrent writes;
- duplicate or missed records after migration;
- partial availability;
- global uniqueness where required.
See Data Testing and Reliability Engineering.
DataOps
Partition changes should be managed through repeatable operational procedures with versioned routing configuration, staged migrations and rollback or forward-repair plans.
See DataOps and Data Platform Operations.
Security Boundaries
Partitions can support isolation between tenants or classifications, but physical separation alone is not an access-control system.
Authentication, authorisation and audit remain necessary.
Encryption and Keys
Some architectures use separate encryption keys by tenant, region or partition to reduce blast radius.
Key placement should follow the security and sovereignty model rather than the partition scheme alone.
Privacy
Partitioning can support data minimisation and local processing, but global indexes, backups and routing metadata may still reveal sensitive information.
All derived control-plane data should be included in privacy review.
Cost
Sharding can increase capacity while also increasing operational cost:
- more nodes;
- more replicas;
- routing infrastructure;
- cross-shard coordination;
- monitoring;
- backup complexity;
- rebalancing work;
- operational expertise.
Do not shard before scale or isolation requirements justify the complexity.
When Not to Shard
A single well-designed database can handle substantial workloads. Premature sharding adds distributed-system failure modes before they are necessary.
Scale vertically, optimise queries, archive old data, add indexes or use read replicas before assuming sharding is the first answer.
Partitioning vs Data Mesh
Partitioning is a physical or logical storage design. Data mesh is an organisational ownership model.
A domain can own one product spread across many shards; one shard does not automatically equal one domain.
See Data Mesh and Federated Data Ownership.
Partitioning and Caching
Caches can be partitioned independently from source databases. Cache keys and source-shard keys do not always need to match.
Different partition strategies may serve different workloads, but each derived state still needs synchronisation and invalidation.
Partitioning and AI
Large AI retrieval stores may partition vectors, documents or tenants across nodes.
Partition design affects recall if a query searches only some partitions. Approximate nearest-neighbour systems therefore need routing strategies that preserve acceptable retrieval quality while controlling cost.
Source lineage and access scope should remain valid across partitions.
Education Example
An education platform partitions student operational data by organisation ID so most school-specific activity remains local. A very large school grows beyond one shard and receives a sub-partition strategy based on student ID.
National-level analytics runs from a separate analytical store rather than scatter-gathering every operational shard for every dashboard refresh.
Commerce Example
An e-commerce platform initially shards orders by customer ID. This keeps customer histories local but creates hotspots for very large business accounts. The routing layer moves those accounts onto dedicated shards while preserving stable order and customer identities.
Time-Series Example
A telemetry platform partitions by day and hashes by device inside each day. Date pruning supports historical queries while hashing spreads high-volume writes across several physical partitions.
Common Failure Modes
- Shard by instinct: the key is chosen without workload evidence.
- Perfect row balance, terrible traffic balance: hot entities overload one shard.
- Sequential-key hotspot: all new writes hit the same range.
- Shard identity becomes business identity: repartitioning breaks consumers.
- Cross-shard everything: locality benefit disappears.
- Global uniqueness assumed: independent shards create duplicates.
- Rebalance without reconciliation: records are lost or duplicated during movement.
- Partial outage becomes zero: global summaries undercount silently.
- Too many tiny partitions: metadata overhead dominates.
- Premature sharding: distributed complexity arrives before real scale need.
A Partitioning and Sharding Checklist
- What scale or isolation problem requires partitioning?
- What are the dominant read and write patterns?
- Which key preserves useful locality?
- How evenly are data and workload distributed?
- Could one tenant, time range or entity become hot?
- Should range, hash, list or composite partitioning be used?
- How are requests routed?
- How is routing metadata protected and versioned?
- Which queries cross partitions?
- Are global uniqueness and secondary indexes required?
- How are partitions replicated and recovered?
- How will shards split, merge or move?
- How is rebalancing reconciled?
- What happens during partial shard failure?
- Can the partition strategy evolve without changing logical entity identity?
A Maturity Ladder
- Split: large data is divided into manageable pieces.
- Keyed: partition assignment follows explicit stable rules.
- Locality-aware: common operations stay near related data.
- Skew-aware: hotspots and uneven growth are observable.
- Rebalanced: partitions can move or split without losing state.
- Resilient: replication and recovery operate per shard and globally.
- Governed: sovereignty, security and identity survive physical movement.
- Adaptive: workload and cost evidence continuously refine distribution strategy.
The Deeper Principle: Scale Is a Locality Problem
Partitioning works by deciding which data should live together and which operations can be separated. That decision determines whether scale creates independence or coordination overhead.
The best partition scheme does not merely distribute bytes evenly. It distributes work sensibly, preserves logical identity, keeps common operations local and leaves a controlled path for rebalancing when today’s pattern becomes tomorrow’s hotspot.
Data Management Series
- Database Management and Transactional Integrity
- Data Lakes and Lakehouses
- Data Synchronisation and Reconciliation
- Data Backup, Recovery and Resilience
- DataOps and Data Platform Operations
Final idea: partition for the workload you actually have, not the symmetry of the diagram. Distribution becomes trustworthy when routing, locality, skew, rebalancing and recovery are designed together rather than added after the first shard becomes too large.
