Skip to content
KEDBYTE
Site navigation
How Data Works
Chapter
35

Partitioning: Dividing the Work

Part F · Data Across Machines|3,701 words|about 16 min read|Volume F

35.0 What this chapter gives you#

  1. Partitioning divides a collection into pieces. It can make some operations smaller, distribute capacity or align storage with the way records are managed. It can also make previously local questions cross several boundaries.
  2. You will choose a partition key from a workload, trace range and hash placement, identify hot partitions and reason about movement, cross-partition queries and identity constraints.
  3. The routing examples are small deterministic models. PostgreSQL table partitioning is documented separately from multi-host sharding; neither a Python routing function nor a partitioned table definition is a distributed database implementation.

35.1 Partition keys#

35.1.1 PLAIN — in simple words#

  1. A partition key is the information used to decide which piece should hold a record. A date can place an order in a month; a branch identifier can place it with other records from that branch.
  2. The choice should match the work. Keeping an order and its lines together can simplify order processing, while grouping by month can simplify bounded historical scans and retention.
  3. No key makes every possible question local. A useful design chooses which operations deserve locality and records the cost imposed on the others.

35.1.2 PLAIN — a picture in your head#

  1. Mira can file invoices by month, branch or customer. Finding one month’s invoices is easy in the first arrangement and may require opening many folders in the others.
  2. Duplicating every invoice into every possible filing scheme creates maintenance and consistency work. The filing decision is therefore about priorities, not discovering a universally correct alphabet.
  3. Where the comparison breaks: databases can maintain secondary indexes and distributed routing metadata. Physical placement and logical access paths need not be identical.

35.1.3 PLAIN — a worked example#

  1. Assume a fictional multi-branch workload where 85% of requests read or update one order within a known branch, 10% summarize one branch and 5% summarize all branches.
  2. A branch-based key may keep much of the first 95% local. The all-branch report still needs a separate aggregation path or access to several partitions.
  3. If one branch later produces 80% of all traffic, the same placement can become badly imbalanced. A workload percentage is a dated assumption to measure, not a permanent property of a branch name.
  4. Record both record volume and request volume. Ten equally sized folders can still receive very unequal amounts of attention.

35.1.4 PLAIN — what is really happening inside#

  1. The application or engine maps the key to a partition. Routing must agree with the current layout version; otherwise a request can reach the wrong owner or miss a moved record.
  2. Related constraints matter. A uniqueness rule spanning all records may become harder if each partition can independently accept the same value.
  3. A partition key that changes during normal work can require moving records. That movement is a state transition with its own atomicity, routing and recovery requirements.

35.1.5 TECHNICAL — the engineer’s version#

  1. Distinguish logical partitioning, physical placement and sharding across independent nodes. PostgreSQL declarative table partitioning divides a logical table into child relations; it does not automatically create a multi-host consensus or transaction layer. [S152]
  2. Review query locality, write locality, cardinality, skew, key mutability and cross-partition invariants. These are workload properties, not consequences of choosing an integer key.
  3. Complete identity must survive placement changes. A record should not acquire a new business identity merely because the partition directory assigns it to another host.

35.1.6 WORDS — remember these#

  1. Partition key: the value deciding placement — a field or expression mapping a record into a defined partition. Locality: related work stays close together — placement or access behaviour reducing the number of independent resources an operation must involve. Shard: an independently placed portion of data — commonly a partition hosted on a separate storage node or database, with product-specific semantics.

35.2 Range and hash approaches#

35.2.1 PLAIN — in simple words#

  1. Range partitioning assigns intervals of ordered values to pieces. Hash partitioning transforms a key into a value used to select a piece.
  2. Ranges can make date or identifier intervals easy to locate. Hashing can spread many distinct keys without preserving their original order.
  3. Both methods need exact boundary rules. A record at a boundary must not disappear between two pieces or be treated as belonging to both.

35.2.2 PLAIN — a picture in your head#

  1. Range filing puts invoice numbers 1–99 in one drawer and 100–199 in another. Hash filing uses a repeatable sorting rule that may place neighbouring invoice numbers far apart.
  2. The first helps with “all invoices from this interval.” The second can spread separate lookups across drawers, but an interval question may need many of them.
  3. Where the comparison breaks: a real hash function has defined encoding and distribution properties. The last digit of an identifier is not automatically a suitable general-purpose distribution scheme.

35.2.3 PLAIN — a worked example#

  1. Define ranges [0,100), [100,200) and [200,300). The left bracket includes the lower bound; the right parenthesis excludes the upper bound. Value 100 belongs only to the second range.
  2. For a deliberately simple hash-routing example, place integer key k into k mod 3. Keys 0, 3, 6 and 9 go to partition 0; keys 1, 4, 7 and 10 go to 1; keys 2, 5, 8 and 11 go to 2.
  3. Now change the rule to k mod 4. Of keys 0–11, only 0, 1 and 2 retain the same numerical destination. Nine of twelve move: 75% in this particular fixture.
  4. This shows why changing the divisor is not a harmless configuration edit. The example’s remainder function is pedagogical, not the PostgreSQL hash algorithm.

35.2.4 PLAIN — what is really happening inside#

  1. Range boundaries need coverage and non-overlap checks. Missing or default partitions require deliberate treatment, especially for new dates or previously unseen tenants.
  2. Hash placement requires stable input encoding, algorithm and layout metadata. A process-randomized language hash must not silently define persistent cross-process routing.
  3. Consistent-hash schemes and indirection through many small logical buckets can reduce how much placement changes when capacity is added. They do not eliminate the need to transfer records safely.

35.2.5 TECHNICAL — the engineer’s version#

  1. PostgreSQL range partitions use inclusive lower and exclusive upper bounds; supported declarative methods also include list and hash partitioning. Their exact operator and constraint semantics are engine-specific. [S152]
  2. The original Dynamo paper describes consistent hashing with virtual nodes. It is a historical design reference, not a specification for current DynamoDB or a universal recommendation for every workload. [S162]
  3. A routing test should cover boundary values, missing keys, encoding changes and layout-version mismatches. Successful placement of a few ordinary keys does not establish safe rebalancing.

35.2.6 WORDS — remember these#

  1. Range partition: a piece owning an ordered interval — a partition defined by lower and upper key bounds. Hash partition: a piece selected through a repeatable mapping — placement based on a hash-derived value rather than original key order. Layout version: the identity of a placement map — metadata distinguishing successive assignments of logical partitions to owners.

35.3 Hot partitions#

35.3.1 PLAIN — in simple words#

  1. A hot partition receives much more work than the others. The system can have spare capacity overall while that one piece is overloaded.
  2. Evenly distributing record counts does not necessarily distribute reads, writes or expensive operations evenly. One popular product can receive more requests than thousands of quiet products combined.
  3. Adding more partitions helps only if the hot work can actually be divided or served differently without breaking its correctness rules.

35.3.2 PLAIN — a picture in your head#

  1. Ten service counters each hold the same number of customer files, but nearly everyone needs the file kept at counter 1. Nine idle counters do not make that queue disappear.
  2. Copying a read-only catalogue description may help. Allowing ten counters independently to promise the same final stock item creates a different problem.
  3. Where the comparison breaks: some systems support coordinated counters, escrow-like allocations or specialized contention management. Those are additional designs, not consequences of drawing more boxes.

35.3.3 PLAIN — a worked example#

  1. Imagine 1,000 requests per second and ten partitions. One key attracts 700 requests per second; all other keys together attract 300.
  2. Even if the remaining traffic distributes perfectly, the hot key’s owner receives at least 700 requests per second while the average is only 100. A dashboard showing average load hides the concentration.
  3. Hashing that same key more carefully does not split its requests: identical keys still map together under a deterministic placement rule.
  4. A possible redesign could separate immutable product-description reads from authoritative stock reservations. Whether that helps depends on which part of the 700 requests is actually expensive and whether the new read path meets its freshness contract.

35.3.4 PLAIN — what is really happening inside#

  1. Measure per-partition and per-key workload where privacy and cardinality limits permit. Request counts, bytes, latency and lock contention reveal different kinds of concentration.
  2. Sequential keys can concentrate recent writes at one range edge. A time-based partition may direct all current events to the newest piece.
  3. Splitting a hot logical entity requires a merge or coordination rule for its state. Without one, the design has merely replaced a bottleneck with inconsistent independent answers.

35.3.5 TECHNICAL — the engineer’s version#

  1. Distribution quality must be assessed against the request distribution, not only the key-space distribution. Hash uniformity cannot remove single-key contention.
  2. Virtual partitions can improve placement flexibility and balancing across nodes, but the original Dynamo discussion also identifies non-uniform load and heterogeneous capacity as design concerns. [S162]
  3. Report maxima and relevant latency distributions alongside averages. A useful intervention is tied to an observed bottleneck and includes a correctness-preserving way to divide or duplicate the work. [S103]

35.3.6 WORDS — remember these#

  1. Hot partition: one piece receives disproportionate work — a localized throughput, latency or resource bottleneck. Skew: an uneven distribution — concentration in key frequency, record size, cost or request rate. Single-key contention: many operations compete on one identity — a workload limitation that ordinary redistribution of unrelated keys cannot remove.

35.4 Rebalancing#

35.4.1 PLAIN — in simple words#

  1. Rebalancing moves responsibility or data between pieces as the system grows or load changes. It is a migration, not simply editing an address.
  2. The process must account for writes that happen while copying, requests using an older map, retries and incomplete transfers.
  3. A successful copy is only one checkpoint. The system also needs a safe authority switch and evidence that the old owner can no longer create an incompatible continuation.

35.4.2 PLAIN — a picture in your head#

  1. Mira moves one filing drawer to a new branch while staff are still adding invoices. Copying the drawer at noon and changing the sign at one o’clock can miss everything added during that hour.
  2. Staff holding an old directory may continue sending papers to the original branch. The move needs a forwarding or rejection rule and a clear cutover boundary.
  3. Where the comparison breaks: physical papers do not normally arrive twice after a timeout. Distributed transfers must handle duplication and stale routing explicitly.

35.4.3 PLAIN — a worked example#

  1. In a toy directory, bucket 12 belongs to A under layout 4. A snapshot is copied to B at sequence 100, but A then accepts changes 101–105.
  2. Switching bucket 12 to B at layout 5 before B has those changes would omit accepted work. The transfer must apply the required suffix or pause and reconcile under a controlled procedure.
  3. At cutover, the design records a valid ownership transition and prevents A from accepting layout-4 writes as current. Stale clients receive a redirect or an explicit retryable routing error according to the protocol.
  4. The copy and suffix examples describe prerequisites. The companion routing model does not implement an online rebalance service or prove a zero-downtime migration.

35.4.4 PLAIN — what is really happening inside#

  1. A bounded pause can simplify a small migration: stop eligible writes, establish the final source boundary, finish the copy, verify, transfer authority and resume. Its availability cost must be declared.
  2. An online design may capture changes during copying, but it then needs ordering, duplicate handling, compatible schemas and a valid final handover. Naive independent dual writes can leave destinations different.
  3. Preserve rollback or roll-forward evidence. Once new authoritative writes occur at B, returning traffic to an old snapshot at A is not automatically a safe rollback.

35.4.5 TECHNICAL — the engineer’s version#

  1. Treat data movement as an ownership protocol with snapshot position, catch-up position, routing version and fencing conditions. Atomicity of a local copy command does not establish atomic distributed cutover.
  2. Rebalance bandwidth competes with foreground operations, replication and recovery. Rate limits and headroom belong in the plan, alongside correctness checks.
  3. Engine-specific partition attachment and detachment have their own locking and validation rules. PostgreSQL’s declarative partition operations are not equivalent to moving an independently writable shard between hosts. [S152]

35.4.6 WORDS — remember these#

  1. Rebalancing: changing where pieces are served — controlled movement of data or ownership to meet capacity and workload needs. Cutover: the point a new path becomes authoritative — a transition with explicit prerequisites and postconditions. Catch-up suffix: changes after an initial copy boundary — the remaining ordered work needed before a destination matches the selected source state.

35.5 Cross-partition questions#

35.5.1 PLAIN — in simple words#

  1. A question involving several pieces needs a rule for combining their answers. Adding subtotals can be valid; averaging averages often is not.
  2. The pieces may also be observed at different times. A report can mix incompatible states even when each local query is correct.
  3. A distributed write is harder still: changing two independent owners does not automatically become one all-or-nothing operation. Chapter 38 develops that boundary.

35.5.2 PLAIN — a picture in your head#

  1. Mira asks each branch for sales totals. She must know whether all branches used the same date range, currency, exclusions and revision boundary.
  2. If one branch reports “today so far” at noon and another at three o’clock, adding their numbers does not create a report for a single common instant.
  3. Where the comparison breaks: distributed systems may support coordinated snapshots or timestamp protocols. These require explicit guarantees; they are not provided by sending requests at nearly the same time.

35.5.3 PLAIN — a worked example#

  1. Branch A has two hypothetical order totals, 100 and 300 paise. Its sum is 400, count 2 and average 200. Branch B has eight orders of 100 paise, giving sum 800, count 8 and average 100.
  2. The overall average is (400 + 800) / (2 + 8) = 120 paise. Averaging the branch averages gives (200 + 100) / 2 = 150, which answers a differently weighted question.
  3. Correct partial aggregation carries the information needed for the final calculation: sum and count here, not only average.
  4. These are separate fictional report values, not replacements for the canonical O-1042 and O-1043 amounts.

35.5.4 PLAIN — what is really happening inside#

  1. A scatter-gather query sends work to relevant partitions and combines responses. The slowest necessary response, retries and network traffic can dominate completion time.
  2. Partial results need explicit labelling. Returning nine of ten partitions without disclosing the missing one can produce a plausible but incomplete business total.
  3. Joins across placement boundaries may require moving rows or partial results. The optimizer and data design must account for bytes transferred as well as local computation.

35.5.5 TECHNICAL — the engineer’s version#

  1. Aggregation must preserve grain, population and mergeable state. Sums and counts can be combined under compatible definitions; averages require their weights, and exact distinct counts cannot generally be added across overlapping populations. [S19]
  2. A local transaction boundary does not imply a globally consistent snapshot across independent databases. Specify the supported distributed-read or reconciliation procedure.
  3. For partitioned PostgreSQL tables, uniqueness constraints must satisfy documented partition-key restrictions. A local unique index on each independently managed shard is not a global uniqueness service. [S152]

35.5.6 WORDS — remember these#

  1. Scatter-gather: ask several pieces and combine their replies — a distributed query pattern with completeness and coordination requirements. Partial aggregation: summarize locally before combining — retaining sufficient intermediate state for the intended final calculation. Global invariant: a rule spanning several owners — a correctness condition not established by independent local checks alone.

35.6 Growth without permanent shortcuts#

35.6.1 PLAIN — in simple words#

  1. A design chosen for today’s scale should make its assumptions visible so they can be revisited. Hard-coding “all important customers are in partition 0” creates hidden dependency rather than a growth strategy.
  2. It is often sensible to begin with fewer moving parts. The requirement is not to distribute everything immediately, but to avoid confusing a temporary placement decision with a permanent business truth.
  3. Stable identities, documented boundaries and measured workloads make later changes easier to reason about.

35.6.2 PLAIN — a picture in your head#

  1. Mira labels folders with durable invoice identities and keeps a separate directory showing where they are stored. A moved folder changes the directory, not the invoice’s identity.
  2. Writing the shelf location into every invoice number would make every reorganization a renaming exercise with many opportunities for mismatch.
  3. Where the comparison breaks: indirection itself has cost and failure modes. A directory must be available, versioned and protected; it is not a free extra layer.

35.6.3 PLAIN — a worked example#

  1. A toy ring has owners at positions 20, 60 and 90 on a 0–99 circle. A key goes to the first owner clockwise, wrapping to 20 after 90.
  2. Key 25 belongs to owner 60. Adding an owner at 40 moves keys in the interval (20,40] to the new owner while leaving the other intervals’ ownership unchanged in this model.
  3. Key 25 now maps to 40; key 65 remains at 90; key 95 remains at 20. This is placement arithmetic, not evidence that the records were transferred or that concurrent requests were handled safely.
  4. The lab tests the deterministic mapping and boundary convention separately from any migration claim.

35.6.4 PLAIN — what is really happening inside#

  1. Separate logical identity, partition identity and physical location. They may change on different schedules and should not be collapsed into one ambiguous field.
  2. Record what would trigger a redesign: sustained hot-partition latency, recovery-time pressure, storage growth or an invariant that no longer fits the existing owner boundary.
  3. Reassess the full workload after a change. Improving write distribution may worsen reports, backups, index maintenance or operational recovery.

35.6.5 TECHNICAL — the engineer’s version#

  1. Consistent hashing provides a placement technique with limited remapping under membership changes, subject to its precise scheme. It does not itself supply replication, conflict resolution, migration atomicity or authorization. [S162]
  2. Maintain versioned placement metadata and test invalid or stale maps. Treat routing configuration as part of the data system’s correctness surface.
  3. Growth decisions should use measured distribution and operational requirements. The companion examples demonstrate mapping and arithmetic; no production scalability or online-rebalance claim follows from them.

35.6.6 WORDS — remember these#

  1. Indirection: use a directory instead of embedding a location — a mapping layer separating stable identity from changing placement. Consistent hashing: a scheme limiting remapping when owners change — a family of placement techniques whose exact balance and movement properties depend on the design. Placement invariant: a rule about where records belong — a condition linking keys, layout versions and valid owners.

35.97 Practice and worked answers#

  1. Which range owns 100 under [0,100) and [100,200)? The second. Half-open intervals avoid overlap at the boundary.
  2. How many keys 0–11 change numerical destinations when mod 3 becomes mod 4? Nine. Only 0, 1 and 2 retain their old destination in this fixture.
  3. Will a uniform hash split 700 requests to one key across ten owners? Not under ordinary deterministic single-owner placement. The same key continues to map together.
  4. What is missing from “copy bucket 12, then change its address”? A boundary for concurrent writes, catch-up, ownership transfer, stale routing, validation and failure handling.
  5. How should averages from branches with different order counts be combined? Combine compatible sums and counts, then divide. Do not take an unweighted average of averages unless that is the intended metric.
  6. May a report hide that one required partition failed to answer? No. Completeness is part of the result contract.
  7. Does a local unique index on every shard ensure a globally unique email address? No. Independent owners can accept the same value unless a broader mechanism or scoped identity rule prevents it.
  8. What does the ring example establish? Its placement and limited-remapping behaviour under the stated boundary rule. It does not establish safe data movement or distributed correctness.

35.98 Common wrong ideas#

  1. Wrong: partitioning always means multiple machines. Right: a table can be partitioned within one database server.
  2. Wrong: equal record counts mean equal load. Right: request frequency and operation cost can be highly skewed.
  3. Wrong: a better hash fixes every hotspot. Right: single-key contention needs a different workload or coordination design.
  4. Wrong: changing the partition count is just a configuration update. Right: it may change record placement and require a migration.
  5. Wrong: copying data transfers authority. Right: authority and stale-writer prevention need explicit cutover rules.
  6. Wrong: local correctness automatically composes globally. Right: cross-partition snapshots, writes and invariants need their own guarantees.
  7. Wrong: partial query answers can be silently treated as complete. Right: missing partitions change the meaning of the result.
  8. Wrong: consistent hashing is a complete database. Right: it is one placement technique among many other required mechanisms.

35.99 Chapter summary in 20 lines#

  1. Partitioning divides a logical collection into pieces.
  2. A partition key should follow measured work and explicit priorities.
  3. Record volume and request volume are different distributions.
  4. Logical partitioning does not automatically imply multi-host sharding.
  5. Range boundaries require complete and non-overlapping rules.
  6. Hash placement requires stable encoding and mapping definitions.
  7. Changing a modulus can move many keys.
  8. Consistent hashing can limit remapping under a specified scheme.
  9. Placement arithmetic does not move or validate records.
  10. A single popular key can overload one otherwise balanced partition.
  11. Monitor per-partition pressure rather than only averages.
  12. Rebalancing is a data and authority migration.
  13. Initial copies need compatible catch-up and cutover boundaries.
  14. Stale routing and old writers need explicit handling.
  15. Cross-partition reports require compatible populations and snapshots.
  16. Combine sufficient aggregate state rather than blindly averaging averages.
  17. Missing required partitions must remain visible in results.
  18. Global invariants are not guaranteed by independent local constraints.
  19. Separate stable business identity from changing physical location.
  20. Revisit placement using workload evidence and operational requirements.

Return to contents