skip to content

Long-Lived State

What a long-running job must remember between records: what it holds, where it is kept, how it stays bounded, and what survives a code change. It is the job's real storage system.

on this pageshow

explore

questions

21

A long-running job holds per-key running totals. Why does deploying new code against it need a plan a stateless service does not?

level: juniorimportance: must knowfreq 60%

answer

  1. a service keeps nothing; this job remembers
  2. the new build inherits bytes
  3. two matches: which step, which format
  4. failure is refusal, or a silent empty start

basics

~20 s

New code has to read back the retained set the previous build wrote — everything the job holds between records. A stateless redeploy discards nothing; here, entries the new build cannot match or decode mean a refused start, or history silently lost.

solid answer

~50 s

A stateless service can be replaced process by process because nothing of value lives inside it. A long-running job's **retained set** — the running accumulators, records buffered awaiting a group, deduplication entries, join sides and pending per-key callbacks it holds between records — was written by the build you are replacing, so the new build has to find it and decode it. Two independent matches have to succeed: the runtime must decide which stored entries belong to which stateful step of the new program, and the new code must be able to read the stored values in the format the old code wrote them. Runtimes differ in what happens when either fails — some refuse to start and name the step, some start that step with an empty retained set and resume from nothing. A batch model that keeps nothing between runs has no such problem: you just run the new code.

go deeper

for a junior

Recall that a job holding running totals, deduplication entries or join sides carries data the next build must read back, and that a stateless service carries none. Naming that difference is the whole of the junior answer.

for a middle

Explain the two matches that must succeed — which stateful step the stored entries belong to, and whether the new code can decode them — and say that a failure may be a refused start on one runtime and a silent empty start on another.

for a senior

Show that you verify a history-dependent number after the deployment rather than trusting a green job, and that you know the fallback's cost grows with the age of the retained set and depends on the source still holding that period.

for a principal

The call you own is the standing convention: what every long-lived job in the organisation must decide before its first deployment about step identity, stored formats and expiry, so that no team discovers in week three that its only route is to start again.

## What the retained set is, and why it is not a cache A long-running job's **retained set** is everything it is still holding between records: running accumulators, records buffered awaiting their group, deduplication entries, both sides of a join, and pending per-key callbacks. It is the job's real storage system, and it is not a cache. Drop a cache and the next answer is slower; drop the retained set and the next answer is **wrong** — a running total restarts at zero, a deduplication entry no longer suppresses the duplicate that arrives an hour later, one half of a join never finds its partner. That is the entire reason replacing the code of such a job is a different operation from replacing a stateless service. A stateless service keeps nothing worth keeping in the process, so a new build can be started, given traffic and the old build killed; nothing is read back because nothing was ever written. A stateful job's new build inherits bytes written by the build it replaces. ## The two matches a redeploy has to make | Match | What is compared | What failure looks like | |---|---|---| | **Identity** | each stateful step of the new program against the label the stored entries carry | entries belong to no step the new program has, or a step starts and finds none of its own | | **Format** | the new code's idea of a stored value against the bytes the previous build actually wrote | a decode error, or — worse — a successful decode into a value that is quietly not what was stored | Both have to succeed, and they fail independently. Code that changed nothing about what is stored can still lose its entries because the step it belongs to is no longer identified the same way; a step that is identified perfectly can still be handed bytes it cannot read. ## What actually happens on a failed match — this varies This is where engineers who have only operated one runtime get caught, because the behaviours are genuinely different: - some runtimes **refuse to start**, naming the step whose entries could not be matched; - some **start that step with an empty retained set**, which is the dangerous one: the job is green, the dashboards are live, and the numbers are quietly missing everything accumulated before the deployment; - some start and **fail later**, at the first read of an entry the new code cannot decode, which can be hours after the deployment looked successful; - some require an explicit acknowledgement before they will proceed with entries they could not place. The operational rule that follows is the same on all of them: after a redeploy, check a number that depends on accumulated history, not just that the job is running. ## Where the problem does not exist at all 1. **The two-phase disk-to-disk batch model** keeps nothing between runs — every run re-derives its result from the input — so a new build is simply the next run's code. Long-lived state questions do not apply to it. 2. **Finite jobs generally**, where the whole answer is computed from an input that is read again each time. 3. **Stateless steps** inside an otherwise stateful continuous job: a filter, a projection, or an enrichment that looks a value up outside the job holds nothing between records, so it can be changed freely. Saying which of these you are in is the first move, not a detail: a great deal of anxiety about deployments is spent on jobs that hold nothing. ## Why the plan is a first-deployment decision The identity a stored entry carries is written **with the entry**, by whatever rule was in force at the time. That makes two things irreversible in practice: - an explicit name pinned to a stateful step after the fact labels the step, not the entries already written under the old label, so the new build looks for a name that nothing on storage carries; - the older the job, the more expensive the fallback, because rebuilding the retained set from the source is only possible for as long as the source still holds the period the entries cover, and it costs a catch-up run over all of it. So the compatibility plan — how a stateful step will be identified, and what happens to a stored value when its shape changes — belongs in the first deployment of a job that is expected to live, alongside the expiry rule that keeps the retained set bounded. It is cheap to decide on day one and can be impossible to decide in week three.

  • The job had been running for six weeks before the redeploy. Does the age of the retained set change anything?
    It changes the fallback, not the mechanism. The older the job, the more history the entries represent, so rebuilding them from the source costs a catch-up run over that whole period and is possible only while the source still holds it. Age also makes a silent empty start harder to notice, because the wrong numbers look plausible for a while.
  • Does a nightly job that re-reads its whole input and rewrites its output have this problem?
    No. It keeps nothing between runs, so the new build simply runs. The distinction is not batch against continuous, it is whether anything is carried from one run or one record to the next: a finite job that persists nothing is free to change, and a continuous job of only stateless steps very nearly is.
  • Is it enough that the new build compiles against the same value class the old one used?
    No. What is on storage is bytes, not a class. The code that turns a stored value into bytes and back can change — a different encoder, or a different configuration of one — while the class in the source looks identical, and then the format match fails even though the code matches.

Replacing a stateless service is swapping a vending machine for a newer one. Replacing a stateful job is swapping the machine while keeping the old one's coin box — the new machine has to open that box and count what is in it, and if the lock does not fit you either stop or start the day's takings at zero.

saying these in an interview costs you the question

  • Calls a stateful job's redeploy just a rolling restart
  • Treats the retained set as a cache that can be dropped without changing results
  • Assumes every runtime refuses to start rather than resuming with empty state
  • Assumes rebuild-from-source is always available as a fallback
  • Says the deployment tooling handles it, naming nothing that reads the stored values
  • Checks only that the job is running after the deployment, never a history-dependent number
open as a page

In a job that keeps a running total per customer, why can a stateful step read only the current record's customer total?

level: juniorimportance: must knowfreq 72%

basics

~20 s

Key-bound state is addressable only under the grouping key of the record in hand. Each key's entry lives with the worker that owns that key, and the step is handed no way to reach another key's entry.

open as a page

A job deduplicates events by holding every id it has seen. Why does that retained set grow forever, and what bounds it?

level: juniorimportance: must knowfreq 70%

basics

~20 s

A deduplication check only ever adds: each new id creates an entry and nothing deletes one, so the retained set grows with distinct ids for the life of the job. Only a removal rule bounds it.

open as a page

A continuous job runs a filter, a per-user running total and a repeat-identifier dropper — which of the three hold something between records?

level: juniorimportance: must knowfreq 70%

basics

~20 s

The running total and the repeat-identifier dropper hold entries between records: one accumulator per user, one entry per identifier already seen. The filter holds nothing, because its verdict depends only on the record in hand.

open as a page

A job keeps everything it remembers between records as live objects inside the worker process - what caps that, and what happens at the cap?

level: juniorimportance: must knowfreq 58%

basics

~20 s

The cap is the memory of that one worker process, shared with everything else the process is doing. There is normally no overflow path, so passing it degrades and then kills the worker rather than quietly moving entries to disk.

open as a page

When a job is redeployed, what matches each stateful step to the entries the previous build wrote for it?

level: middleimportance: must knowfreq 55%

basics

~20 s

An identity label written with the entries. Either the author pinned a stable step identifier — a name fixed to that stateful step — or the runtime derived one from the step's position in the job's step graph, which later edits can shift.

open as a page

Why does a step with key-bound state require a redistribution by key upstream, and what does that routing buy?

level: middleimportance: must knowfreq 64%

basics

~20 s

Exactly one worker handles every record for a given key, so its entry has a single writer and needs no lock, no network call and no consistency protocol to read. The upstream redistribution is what gets records there.

open as a page

An expiry rule can measure its interval from an entry's last write or its last read. How do the two differ, and which suits deduplication?

level: middleimportance: must knowfreq 56%

basics

~20 s

Last-write expiry gives an entry a fixed lifetime from its last store or update; last-access expiry restarts the clock on reads too, so hot entries never leave. Deduplication wants last-write: a read there means a duplicate arrived.

open as a page

A continuous job ingests only two thousand records a second yet retains a few hundred gigabytes — which quantities actually set that size?

level: middleimportance: must knowfreq 62%

basics

~20 s

Distinct live keys, bytes per entry and per-entry overhead set the size, with the retention horizon deciding which keys count as live. Input rate never enters for accumulators or deduplication entries, so a slow job can retain hundreds of gigabytes.

open as a page

What changes when a stateful step moves its retained set from process memory into an embedded on-disk store with a memory cache?

level: middleimportance: must knowfreq 66%

basics

~20 s

The ceiling moves from one worker process's memory to its local disk, and in exchange every read and write pays an encode and a decode, a disk read whenever the cache misses, and the background merging of the store's own files.

open as a page

Which changes to a stateful step's stored values can still be read back after a redeploy, and which cannot?

level: middleimportance: should knowfreq 50%

basics

~20 s

Additive changes usually survive: a field added with a default supplied for entries that lack it, and a widened type. Renames, a changed grouping key, a swapped record encoder and removed or reordered stateful steps do not, because the stored bytes no longer answer what the new code asks.

open as a page

A source step tracks its own read position rather than anything per key; how is that state redistributed when the width changes?

level: middleimportance: should knowfreq 46%

basics

~20 s

That is per-worker state: attached to one parallel instance of a step, not to a key. With no bucket to move, it is redistributed by a rule the author picks — the collected entries split across new instances, or every instance given the whole set.

open as a page

A job clears each key's entry with a per-key scheduled callback. How does that bound the retained set, and how can the callbacks themselves leak?

level: middleimportance: should knowfreq 44%

basics

~20 s

The callback is a wake-up registered against one key and one moment; what bounds the retained set is the handler that deletes the entry when it fires. Pending callbacks are stored entries too, so one per record leaks just as badly.

open as a page

A job holding months of per-key entries must change its grouping key. What are the options, and what does each cost?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Three routes: rewrite the stored snapshot offline into the new key and start from it; run a transition window where the new build writes the new key and still reads the old; or start empty and rebuild from the source. Each trades downtime, complexity or a period of wrong numbers.

open as a page

A stateful continuous job started at four workers must now run at forty; why can the key buckets fixed at its first start block that?

level: seniorimportance: should knowfreq 53%

basics

~20 s

The key space is cut into a fixed number of buckets at the job's first start, and a key never changes bucket. A width change therefore moves whole buckets between workers, and can never use more workers than there are buckets.

open as a page

Your team is adding a removal rule to a long-running job's retained set. What correctness does each rule cost, and what decay is it buying off?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Every removal rule trades correctness for survival: dropped deduplication entries readmit duplicates, a dropped join side silently loses matches, a cleared activity entry splits one burst in two. It buys a set that stops driving snapshot size and restart time upward.

open as a page

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%

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.

open as a page

A team sized its retained set by counting per-key accumulators alone and came out tenfold low — which held entries did that count miss?

level: seniorimportance: should knowfreq 48%

basics

~20 s

It missed records buffered awaiting their group, both sides of a join held until a match is possible, deduplication entries, pending per-key wake-ups, and the per-entry overhead the holder adds. Each is counted on a different quantity from accumulators.

open as a page

A worker retains 40 GB of state but touches only 2 GB per cycle and has 16 GB of memory - which placement fits, and what could overturn it?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Forty gigabytes cannot live as objects in a sixteen-gigabyte process, so a disk-backed store with a memory cache fits, and the two gigabytes touched per cycle are mostly cache hits. Restore time or an unstable hot set can overturn it.

open as a page

Across dozens of long-running jobs, what standing rule would make each team predict its retained set before first deployment?

level: principalimportance: should knowfreq 38%

basics

~20 s

Require a written retained-set budget per job before first deployment: entry categories, live key cardinality with its source, the horizon per category, bytes per entry, and the product. Fix the units so numbers compare, and require re-declaration when the grouping key or horizon changes.

open as a page

A continuous job runs as repeated small finite jobs and writes changed state entries to a versioned file set each cycle - what dominates its state cost?

level: seniorimportance: nice to knowfreq 32%

basics

~10 s

Cost is paid per cycle, not per access: each cycle writes a file per state partition to durable shared storage, plus periodic full copies, so partitions multiplied by cycle rate sets the bill.

open as a page