skip to content

A job running four months now takes an hour to restart instead of two minutes, and per-record latency has doubled. How does an unbounded retained set cause both?

level: seniorimportance: should knowfreq 48%

answer

  1. one cause, two symptoms
  2. restart reads it all back
  3. working set outgrows the cache
  4. slope depends on the store model
  5. trend follows calendar, not traffic

basics

~20 s

Restart time tracks the bytes a job reads back before its first record, and the retained set is those bytes. Latency rises because the working set outgrows the memory the store had for it. Both curves follow one growing entry count.

solid answer

~50 s

Both symptoms are the same quantity seen from two sides. A restart reads the durable snapshot — the consistent copy of the retained set written to storage outside the workers — back before the job processes anything, so restart time follows the bytes that snapshot carries. Per-record latency rises because the part of the set actually touched each cycle stops fitting the memory the store had for it: an in-process memory store pushes the worker's memory management harder as the live object graph grows, an embedded on-disk store with a memory cache starts missing its cache and pays an encode, a decode and a local file read per access while its background file merge takes more of the disk, and state rewritten as files each cycle writes a larger file set every cycle. Nothing broke in week sixteen; the job crossed a line it had approached since week one.

go deeper

for a junior

Know that everything the job remembers has to be written to a durable copy and read back on restart, so a set that keeps growing makes restarts keep getting slower.

for a middle

Explain both curves from one cause, and say how the per-access cost grows differently for entries held as live objects, entries in an on-disk store with a memory cache, and state rewritten as files each cycle.

for a senior

Distinguish the calendar-driven trend from a traffic-driven one, and know that adding a rule stops the curve without rewinding it — recovering space may need an offline rewrite or a rebuild from source.

for a principal

Treat restart time as a first-class service property with a budget: a job whose restore no longer fits the operational window is out of contract, whatever its steady-state throughput looks like.

## Why the two symptoms are one cause **The retained set** is everything the job is still holding between records: accumulators, buffered records, deduplication entries, join sides and pending per-key wake-ups. If nothing removes entries, that set grows with distinct keys for the life of the job. Every one of its bytes is charged twice — once to the periodic **durable snapshot**, the consistent copy written to storage outside the workers, and once to the store the workers read on every access. Restart time is the first charge; per-record latency is the second. The defining property is that neither charge announces itself. There is no failure, no error rate, no threshold crossed — just a job a little slower every week and a restart a little longer, until an ordinary deployment does not finish inside its window. ## Why coming back takes longer - A restarting worker must have its entries in place **before** it processes a record, so the read is on the critical path. - The bytes it reads back are the bytes the retained set holds, which is why restart time follows the same curve as entry count. How much has to be read varies with the store model: where a snapshot is written whole each cycle, it is the full set; where snapshots are incremental, a restore reads a base plus the increments since, which is smaller per cycle but still grows with the base. - More of it is fetched from shared storage across the network at once, by every worker simultaneously — a different bottleneck from the steady state, and why restart time can degrade faster than the set grows. - Where a job is stopped deliberately, the **planned snapshot** taken on the way down has the same size problem: the graceful stop that used to take a minute now takes long enough that people start killing it instead, which is how a clean shutdown turns into a recovery. ## Why each record costs more, by store model | store model | what grows | what the record pays | |---|---|---| | **in-process memory store** — entries kept as ordinary live objects inside the worker | the live object graph in the worker process's managed memory | more time lost to the process's own memory management, and eventually a hard ceiling at that process's limit | | **embedded on-disk store with a memory cache** — recent entries in memory, the rest in local files merged in the background | the share of accesses that miss the cache | an encode, a decode and a local file read per miss, plus a background file merge competing for the same disk | | **state rewritten as files each cycle** — each cycle writes changed entries to a versioned file set, with a periodic full rewrite | the size of each cycle's file set and of the periodic full copy | a larger fixed cost per cycle, paid whether or not many records arrived in it | The first model degrades suddenly and the second gradually; the third degrades in steps at each full rewrite. Naming which model you mean is the difference between an answer that is true and one true only of the engine you used last. ## Why it only appeared in month four - The growth is linear in distinct keys, but the **cost curves are not**: an access is flat until the working set exceeds the cache, then steps up; a restart is fine until it exceeds the operational window, then it is an outage. - Tests never see it. A test runs for seconds over a handful of keys, and the defect is defined by months and millions of keys. - The early symptoms are attributed elsewhere — a slow week is blamed on traffic, a slow restart on the storage layer — because nothing in the job changed. ## Confirming it is the retained set and not something else The distinguishing evidence is that the trend is monotonic across restarts and independent of input rate: latency does not fall when traffic falls, and each restart is slower than the one before it even when the intervening period was quiet. A slowdown that tracks traffic is a capacity story; a slowdown that tracks *calendar time since the job started* is this one. Comparing the size of successive snapshots over weeks settles it, and which numbers a team watches routinely is another node's subject. ## Getting out of it 1. Add the removal rule. It stops the curve but does not by itself rewind it. 2. Expect the bytes back later than the entries. Some stores remove eagerly; others mark the entry and reclaim space on its next touch or during a background file merge, so a set of entries never touched again can sit on disk long after they logically expired. 3. If space does not come back in time, the remaining options are to rewrite the stored snapshot outside the running job — reading it as an ordinary dataset, dropping what should have expired and writing a new one the next run reads — or to discard the retained set and rebuild from the source, which costs exactly the facts you were keeping. 4. The two-phase disk-to-disk batch model never has this problem, since it keeps nothing between runs: a continuous job inherits a maintenance obligation its batch predecessor did not have.

  • The team adds an expiry rule today. When does the restart get faster again?
    Later than you expect. The entry count falls as the rule takes effect, but bytes come back only as fast as the store reclaims them: eager removal frees space immediately, while lazy removal reclaims on the next touch or during a background file merge, which may never visit an entry nobody reads. If the deadline is urgent, plan on rewriting the stored snapshot offline or rebuilding from source.
  • How would you separate this from a slowdown caused by rising traffic?
    By whether the trend follows the calendar or the load. Unbounded growth degrades monotonically across restarts and does not improve during quiet periods, and each successive snapshot is larger than the last. A capacity problem tracks input rate in both directions. The snapshot size series over weeks is the cheapest discriminator.
  • Why can a graceful stop become dangerous as the retained set grows?
    Because the stop writes a snapshot on the way down, and that write grows with the set. Once a clean shutdown routinely exceeds people's patience they start killing the job instead, which turns an orderly stop into a recovery from the last periodic copy — replaying more input and, where the output is not idempotent, re-emitting effects the world already saw.

saying these in an interview costs you the question

  • Blames the storage layer because nothing in the job's code changed.
  • Expects an out-of-memory failure first rather than a slow restart.
  • Assumes adding expiry frees the bytes immediately in every store.
  • Thinks more workers fix it, when the total set is unchanged.
  • Describes one store model's degradation as if all engines behaved that way.
  • Reads a slow week as traffic, without checking snapshot size over time.