skip to content

In a send-later service storing billions of future jobs, how do time-bucket partitions and sharded pollers let due-job polling scale horizontally?

level: seniorimportance: should knowfreq 45%

answer

  1. one hot 'now' range
  2. time slice plus hash
  3. bucket acts as the index
  4. per-shard drained checkpoint
  5. shard count recorded per bucket

basics

~20 s

Partition jobs by due-time bucket and shard number, hashing each job to a shard. Each poller owns some shards and reads only the current bucket for them, so adding shards and pollers spreads both writes and polling.

solid answer

~50 s

A single due-time index eventually becomes one hot range: every write for a popular time and every poll hit its 'now' end. Instead, key jobs by `(bucket, shard)`, where `bucket` is the due time truncated to, say, a minute and `shard = hash(job_id) mod S`. Writes for a popular minute spread across `S` partitions, and the bucket itself acts as the index, so even a store without secondary indexes works. An assignment mechanism gives each poller a subset of shards; each poller reads the current bucket for its shards, dispatches, and records a per-shard checkpoint of the last fully drained bucket, so a restart resumes where it stopped and lagging buckets are caught up. A job scheduled into a bucket that has already passed is written into the current one. Record `S` per bucket so the shard count can grow for future buckets.

go deeper

for a junior

Remember the key shape: a time slice plus a hash-based shard number, so no single partition holds all the jobs due at one moment.

for a middle

Explain how pollers own shards, read the current bucket, and why the bucket replaces a global index in stores without secondary indexes.

for a senior

Cover checkpoints, catch-up after downtime, late writes into closed buckets, lag metrics, and changing shard count safely with a cutover.

for a principal

Size it: pick the bucket width and shard count from peak due rates, and decide how shard assignment and rebalancing are operated.

## Why one index runs hot A **send-later service** stores future **one-shot jobs** and fires each at its **due time**. The first design is one table with an index on `due_at`, polled for `due_at <= now`. At billions of jobs and several thousand due per second, that design concentrates load: - Every **poll** reads the same end of the index — the 'now' end. - Popular times (09:00, midnight) put millions of **writes** into one narrow key range. - Many large stores, especially **wide-column** and **key-value** stores, have weak or no secondary indexes, so 'sorted by due time across the whole dataset' may not exist at all. The fix is to make time part of the **partition key** and to split each time slice further. ## Bucket-and-shard keys - **Bucket:** the due time truncated to a fixed width, such as one minute: a job due at 09:00:37 belongs to bucket `09:00`. - **Shard:** a number from `0` to `S − 1`, computed as `hash(job_id) mod S`, so jobs are spread evenly regardless of their due time. - **Partition key:** `(bucket, shard)`. Inside a partition, rows can be sorted by exact `due_at` for finer precision. The bucket plays the role of the index: to find jobs due now, you read partitions whose bucket is the current minute. No global sort is needed. ```json { "partition": { "bucket": "2027-03-01T09:00Z", "shard": 17 }, "due_at": "2027-03-01T09:00:37Z", "job_id": "j-81f2", "status": "pending" } ``` ## The poller loop 1. An **assignment mechanism** (a membership service, consistent hashing over live pollers, or a static map) gives each poller a set of shards. 2. For each owned shard, the poller reads its **checkpoint**: the last bucket it fully drained. 3. It reads the next bucket after the checkpoint, up to the current one, and claims the rows that are due. 4. It dispatches the claimed jobs. 5. Once a bucket is in the past and fully drained, it advances the checkpoint. Pollers with different shards never read the same partition, so they scale out almost linearly until the store or the workers become the limit. ## Checkpoints and late writes - **Restarts and lag.** Without a checkpoint, a poller that was down for ten minutes would jump to 'now' and silently skip ten buckets. With one, it drains the missed buckets oldest-first. - **Lag is a metric.** `now − checkpoint` per shard shows exactly which shards are falling behind. - **Late writes.** A job created at 09:00:50 for 09:00:10 belongs to a bucket that may already be drained. Write it into the *current* bucket, or into an 'immediate' lane, never into a closed bucket. - **Keep the current bucket open.** The poller keeps re-reading the current bucket until it is in the past, because new jobs for this minute can still arrive. ## Sizing Assume, illustratively, a peak of 1,000,000 jobs due in one minute and `S = 100`: - Each partition holds about 1,000,000 / 100 = **10,000 jobs**. - Each partition must be drained at about 10,000 / 60 ≈ **167 jobs per second**. - With 10 pollers each owning 10 shards, each poller handles about **1,667 jobs per second**. The bucket width is a trade-off: | Bucket width | Partition size | Empty reads | Precision inside bucket | |---|---|---|---| | 1 second | small | many when traffic is low | exact enough without sorting | | 1 minute | moderate | few | sort by `due_at` inside | | 1 hour | large | very few | large partitions, costly reads | ## Changing the shard count If you change `S` and pollers start using the new modulus, jobs already written with the old `S` land in shards nobody reads. The safe approach is to **record `S` with each bucket** (or in a small schedule table) and apply a new value only to buckets after a cutover time. Old buckets are still read with the old `S` until they drain. Rebalancing pollers across shards is easier: move shard ownership, and rely on checkpoints plus atomic claims so the handover neither skips nor duplicates work. Adding more pollers than shards gives no extra throughput, since there is nothing left to divide — so pick an `S` larger than the poller count you expect.

  • How do you choose the bucket width?
    Balance partition size against empty reads. Narrow buckets keep partitions small and precision high but mean many reads of empty buckets at quiet times; wide buckets mean fewer reads but large partitions that must be sorted by exact due time. A width matching the service's precision contract, often a minute or less, is a common starting point.
  • What happens if one poller falls several buckets behind?
    Its checkpoint lags, so it drains the missed buckets oldest-first before reaching the current one, and those jobs fire late rather than never. Watch `now − checkpoint` per shard and alert on it. If lag persists, move some of its shards to other pollers or add shards for future buckets.
  • How do you raise the shard count without losing jobs?
    Store the shard count alongside each bucket and apply the new count only to buckets after a cutover time. Pollers read each bucket with the count it was written with, so jobs written under the old count are still found. Once the old buckets have drained, the old count no longer matters.

saying these in an interview costs you the question

  • Partitioning by due minute alone is enough; no shard number is needed.
  • Pollers can always read just 'now' and never need a checkpoint.
  • Changing the shard count only means pollers start using the new modulus.
  • A job created late should go into its already-drained original bucket.
  • Running more pollers than shards keeps raising throughput.