skip to content

How do you update a running Dataflow streaming job in place without losing its state?

level: seniorimportance: must knowfreq 60%

answer

  1. you do not stop the old job first
  2. the job name is the identity, not the ID
  3. the service checks the graph before swapping
  4. renamed steps need a mapping to keep state
  5. a failed check leaves production untouched

basics

~20 s

Submit the new pipeline with the --update flag and the running job's name. Dataflow runs a job-graph compatibility check; if it passes, the replacement job takes over the old job's in-flight state and source position. If it fails, the original keeps running.

solid answer

~50 s

You relaunch the same pipeline with `--update` and the running job's name (`--jobName` in Java, `--job_name` in Python). Dataflow compares the new job graph against the running one in a **compatibility check**: it verifies that the intermediate state and buffered data held by the old job can still be interpreted by the new one. If steps were renamed, you supply `--transformNameMapping` / `--transform_name_mapping` as a JSON map from old step names to new ones so the state follows the rename. On success the old job stops, the replacement starts from the old job's state and source position, and the job ID changes while the job name stays. On failure the replacement never starts and the original job keeps running untouched — which is why update is safer than drain-and-restart. Changes that alter how state is encoded, such as a different windowing strategy or a changed coder on a stateful step, typically fail the check.

code

bash · 7 lines
bash
python pipeline.py \
  --runner=DataflowRunner \
  --project=my-project --region=us-central1 \
  --streaming \
  --job_name=events-enricher \
  --update \
  --transform_name_mapping='{"CountByUser":"CountByUserId"}'

go deeper

for a junior

Know that a running Dataflow streaming job can be replaced in place with a --update launch rather than being stopped and relaunched, and that the job name identifies which job is being replaced.

for a middle

Explain the compatibility check — that Dataflow matches state to steps by transform name and refuses the swap if the buffered state cannot be carried over — and what transform name mapping is for.

for a senior

Demonstrate deploy judgment: which changes break compatibility, why a rejected update is safer than a drain, and how you cut over when the change genuinely cannot preserve state.

for a principal

Own the release policy for stateful streaming: stable transform naming as a code standard, isolating windowing changes into their own deploy, sink idempotency as a precondition, and the cost/rollback tradeoff of parallel-run cutovers.

## Why in-place update exists A stateful streaming pipeline cannot simply be restarted. Its value lives in per-key accumulations, open windows and pending timers that took hours or days to build. Stopping and relaunching throws that away and re-reads whatever backlog is left on the source. Dataflow's answer is the **in-place update**: you submit the new version of the pipeline as a *replacement* for the running job, and the service hands the old job's state to the new one. ## The mechanics You run the same launch command you would normally run, plus two things: the `--update` flag and the **name** of the currently running job. ```bash python pipeline.py \ --runner=DataflowRunner --region=us-central1 \ --streaming --job_name=events-enricher \ --update ``` In Java the equivalents are `--update` and `--jobName=events-enricher`. The job *name* is the identity that matters; the replacement gets a new job ID, and the console shows the old job as "Updated" with a pointer to its successor. ## The compatibility check Before anything is swapped, Dataflow performs a job-graph compatibility check. Its purpose is narrow and worth stating precisely: it confirms that the **intermediate state information and buffered data** held by the running job can be transferred into the replacement. It is not a code review and it does not check business logic. The check is name-based. Dataflow identifies steps by their transform names, so if you renamed a stateful step — or if you changed the code in a way that changed the auto-generated name — the service can no longer match the old state to the new step. That is what `--transform_name_mapping` fixes: ```bash --transform_name_mapping='{"CountByUser": "CountByUserId", "ReadEvents": "ReadRawEvents"}' ``` Each entry maps an old step name to its new name so the state follows. You only need entries for steps that were renamed, and in practice you care most about the stateful ones (aggregations, stateful `DoFn`s, anything holding windows or timers). Changes that commonly **fail** the check, or that are unsafe even when they pass: - Changing the windowing strategy or trigger of a stage that holds window state — the buffered state no longer means the same thing. - Changing the coder used for an element or state type on a stateful step, so old bytes cannot be decoded. - Removing a stateful step whose state was still in use, or replacing it with a differently shaped one. - Renaming steps without a mapping. Changes that are usually fine: adding or removing purely stateless transforms, changing the body of a `DoFn`, adding a new branch or sink, changing a side-input's contents, and adjusting worker machine type or max workers. ## The failure mode is the good news If the compatibility check fails, **the replacement job does not start and the running job keeps running**. You get an error describing the incompatible step. This is the property that makes update the default deployment path: a bad deploy attempt is a non-event for production, whereas a drain-and-restart has already stopped the old job by the time you discover the new one is broken. Always read the failure message rather than reflexively falling back to drain. ## What you do when the update cannot work Some changes genuinely cannot preserve state — a new windowing strategy is the classic example. Then you have real choices, and this is the judgment part of the question: - **Drain and restart.** Drain flushes open windows (partial final results), then you launch the new job fresh. Simple, but the new job re-reads the remaining source backlog and the final windows are truncated, so the sink must tolerate partial and duplicate rows. - **Take a Dataflow snapshot** of the streaming job's state and Pub/Sub source position, then start the new job from the snapshot. This is useful when you want a restore point, though a snapshot restores into a compatible pipeline — it does not rescue an incompatible graph change. - **Run both jobs in parallel** on separate subscriptions, writing to an idempotent or upsert-keyed sink, and cut over when the new job has caught up. This costs double for the overlap but gives you a real rollback. ## Operational habits that make updates boring Name your transforms explicitly and stably — `.apply("CountByUser", ...)` rather than relying on generated names — so refactoring the code does not silently rename steps. Keep windowing strategy changes on their own deploy, separated from routine logic changes. Watch data freshness and system latency for a few minutes after the swap: the replacement briefly re-plays some work, so a small latency bump is normal, and a lasting one is not. And keep the sink idempotent, because both update and its fallbacks assume some records can be written twice.

  • What exactly does --transform_name_mapping do, and when do you need it?
    It is a JSON map from old step names to new step names. Dataflow matches state to steps by name, so if you renamed a stateful transform — or refactored code in a way that changed its generated name — the mapping tells the service which new step inherits the old step's buffered state and timers. Steps you did not rename need no entry.
  • An update fails the compatibility check. What has happened to the production job?
    Nothing. The replacement job is never started and the original keeps running with its state intact. That is the main reason to prefer update over drain-and-restart as the default deploy path: a rejected deploy is a non-event, whereas a drain has already stopped production before you learn the new pipeline is wrong.
  • Which pipeline changes are most likely to make an in-place update impossible?
    Changes that alter how existing state is encoded or interpreted: a different windowing strategy or trigger on a stateful stage, a changed coder for an element or state type, or removing a stateful step. Purely stateless edits — a new DoFn body, an extra branch, a different sink — normally update cleanly.
  • How do you cut over when the change genuinely cannot preserve state?
    Either drain and relaunch, accepting partial final windows and reprocessing of the remaining source backlog, or run the new job in parallel on its own subscription until it catches up and then retire the old one. Both require an idempotent or upsert-keyed sink, since some records will be written twice.

saying these in an interview costs you the question

  • Thinks you must cancel or drain the job before updating
  • Believes the compatibility check validates business logic
  • Assumes any code change updates cleanly regardless of state
  • Says a failed update takes production down
  • Expects the job ID to stay the same across an update

context