skip to content

A year of history must be pushed through a job whose output lands in a live serving database - what sets the replay's duration?

level: seniorimportance: must knowfreq 52%

answer

  1. two capacities in series
  2. the elastic side is not the binding side
  3. rows divided by spare writes per second
  4. cap the concurrent writers yourself
  5. watch replication lag, not cluster health

basics

~20 s

The write rate the destination can sustain while still serving its live traffic, not the cluster. Beyond that point extra workers only produce rejected writes, retries and contention. Plan the replay backwards from the destination's spare capacity.

solid answer

~50 s

Compute is the elastic side and the destination almost never is, so the honest plan starts at the far end. Measure what the destination can absorb while live reads and writes continue: spare write throughput, index maintenance, replication lag, lock contention. Divide the row count by that figure and you have the replay's real duration, which is frequently days rather than hours. Then hold the job to it deliberately - cap how many pieces write concurrently rather than trusting the engine to slow down, size the write batches to what the destination likes, and for a long replay prefer writing files to shared storage and loading them in a controlled pass over row-by-row inserts. Watch destination-side signals as the control: error rate, latency for live traffic, replication lag. The cluster's own numbers will look wonderfully healthy while the database behind it degrades.

code

python · 10 lines
python
rows = 360 * 24 * 3600 * 40          # a year of history at 40 rows/second
spare_writes_per_second = 3000       # measured on the destination WHILE live traffic runs

hours = rows / spare_writes_per_second / 3600
print(round(hours, 1))               # 115.2 -> just under five days

# the cluster is irrelevant to that number; what it controls is only this:
per_writer_rows_per_second = 250
max_concurrent_writers = spare_writes_per_second // per_writer_rows_per_second
print(max_concurrent_writers)        # 12 pieces may write at once, however wide the job reads

go deeper

for a junior

Understand that a job has two speeds - how fast it can compute and how fast its destination will accept the result - and that only the first one grows when machines are added.

for a middle

Do the arithmetic out loud: rows of history divided by the writes per second the destination can spare gives the duration. Explain why retries on a saturated destination make matters worse rather than better.

for a senior

Measure the destination's spare capacity under live load, cap the writing step's concurrency explicitly, stage the output where a bulk load is possible, and alarm on destination latency and replication lag rather than on cluster health.

for a principal

Weigh a five-day throttled replay against the alternatives - a separate destination cut over afterwards, or accepting that some history will never be recomputed - and decide who carries the risk to the live service either way.

## Which limit actually binds A replay has two capacities in series: how fast the cluster can read and transform stored history, and how fast the destination can accept the result. Adding machines moves the first one. It does nothing at all to the second. On any replay whose output lands in a transactional database, a search index, an external service or a broker with downstream consumers, the destination is the binding constraint, and the arithmetic that matters is a division nobody enjoys doing: - rows of history to be written, divided by - writes per second the destination can spare while still serving live traffic, - equals the replay's duration, whatever the cluster costs. The failure of the naive plan is not slowness. It is that the destination degrades while serving someone else's traffic, and the first person to notice is a user of the live system rather than the engineer running the replay. ## Measuring the spare capacity The number you need is not the destination's benchmark maximum. It is what it can give you **in addition** to what it is already doing: - Sustained write throughput under the live read load, not a burst figure from an idle system. - The cost of index and constraint maintenance per row, which frequently dominates the raw insert. - Replication or follower lag, which is the earliest honest sign that a write-heavy replay is outrunning the destination even while writes still succeed. - Lock and contention behaviour on the specific rows the replay touches; a replay that updates hot rows is far more disruptive than one appending cold ones. - For an object store destination, the limit is usually request rate against a prefix, plus the small-file count the replay leaves behind, rather than bytes per second. Take the smallest of these, apply a margin, and treat the result as a budget. ## Holding the job to the budget 1. **Cap the concurrency of the writing step**, not the whole job. The read and transform steps may run as wide as you like; the number of pieces writing at once is what the destination experiences. 2. **Size the write batch to the destination, not to the job.** Larger batches reduce round trips up to the point where they lengthen transactions and raise contention, after which throughput falls while error rates climb. 3. **Prefer a staged path for long replays.** Writing the output as files to shared storage, then loading it in one controlled pass, converts an uncontrolled stream of small writes into a load the destination's operator can schedule, throttle and roll back. 4. **Make writes idempotent.** A replay that must be stopped and resumed, or one whose units of work are retried, will write some rows twice; an idempotent write turns that from a correctness incident into a non-event. 5. **Give the replay a way to pause.** The single most useful operational feature of a multi-day replay is that someone can stop it during an incident and resume it afterwards. ## What varies between engines Do not assume the runtime will protect the destination for you. Whether resistance at the write step propagates back through the graph and slows the read - backpressure, meaning a step unable to hand its output on causes the step before it to stop producing - differs across this family. Where the runtime pushes records step to step as they are produced, that propagation is usually built in and a slow destination does throttle the read. Where a continuous job is assembled from repeated small finite runs, or where the replay is simply a large finite job, there is often no such loop: the write step's units of work slow down, time out, retry and eventually fail the run, having hammered the destination throughout. The cap on concurrent writers is yours to set precisely because you cannot rely on the engine's. ## The signals to watch Read the destination, not the cluster. A replay that is destroying a database shows healthy processed throughput, healthy worker memory and a rising retry count, which is easy to dismiss. The destination shows rising write latency, rising replication lag, rising rejection or timeout rates, and degrading latency for the live traffic that was there first. Alarm on those, and make the replay's write concurrency something an operator can turn down without redeploying the job.

  • The destination is a log the job publishes to rather than a database. Does the constraint disappear?
    It moves. A log will usually absorb the write rate, but every consumer of that log now receives a year of records in a few hours, and the slowest of them becomes the binding constraint instead. You either pace the replay to the slowest consumer, publish to a separate location that only the intended consumers read, or agree with those teams what surge they can take.
  • Why is making the replay's writes idempotent worth the effort up front?
    Because a multi-day replay will be interrupted, and the engine will retry failed units of work on its own. Both replay some rows. If the write is idempotent - keyed and upserted, or written to a location that is replaced wholesale - an interruption costs time only. If it is not, every interruption becomes a correctness investigation and the replay stops being restartable at all.

saying these in an interview costs you the question

  • Doubles the cluster so the replay finishes sooner
  • Assumes the engine's backpressure will protect the destination automatically
  • Treats retried writes as free rather than as extra load
  • Plans the replay from compute cost alone and never measures the destination
  • Believes a larger write batch always raises sustained throughput