How would you size horizontal autoscaling for a Dataflow streaming job with a spiky Pub/Sub backlog?
answer
- start from the SLA and the measured per-worker rate
- the ceiling is a guardrail, not a target
- state on disks makes scaling slow
- the key space caps real parallelism
- the sink has to survive the peak too
basics
~20 sMeasure steady-state throughput per worker, then set the maximum worker count so peak backlog burns down inside your freshness SLA with headroom. Run on Streaming Engine so scaling is not pinned to disks, and confirm key parallelism, sink capacity and quota can actually absorb the extra workers.
solid answer
~50 sStart from measurement, not from a guess. Establish elements per second per worker at a stable input rate, then compute how many workers are needed to clear the worst observed backlog inside the freshness SLA, and set `--max_num_workers` a little above that as a cost guardrail rather than a target. Run the job on Streaming Engine: without it, state lives on worker Persistent Disks provisioned against the maximum worker count at launch, which makes scaling slow and pins your ceiling. Then check that the ceiling is reachable — Dataflow's streaming autoscaler reasons about backlog and CPU utilisation, but real parallelism is bounded by the key space, so a stage keyed on a handful of values will not use the workers you paid for. Verify the sink and any external service can absorb peak worker count without throttling, and watch the autoscaling event log plus data freshness after a spike to confirm the ceiling, not the pipeline, was the constraint.
code
bash · 6 linespython pipeline.py \
--runner=DataflowRunner --region=us-central1 --streaming \
--job_name=events-enricher \
--enable_streaming_engine \
--num_workers=4 \
--max_num_workers=40go deeper
Know that a Dataflow streaming job scales its worker count automatically between a starting size and a configured maximum, and that the maximum is something you set at launch.
Explain the signals the streaming autoscaler reacts to — backlog and CPU utilisation — and why a maximum worker count is an upper bound rather than a guarantee of that much parallelism.
Diagnose a job that will not scale: check key-space parallelism, sink and external-service throttling, quota, and whether state pinned to worker disks is making scaling slow and expensive.
Own the capacity model: derive the ceiling from a freshness SLA and a measured per-worker rate, budget for the spike rather than the average, decide when fixing the pipeline beats buying workers, and hold teams to jobs that actually scale back down.
## Frame the question as capacity planning, not flag-setting The weak answer is "set max workers high and let autoscaling handle it". The real answer starts by naming the objective: a **freshness SLA** — how far behind real time the output is allowed to fall, and for how long — and a **cost envelope**. Everything else is derived from those two numbers. ## Step 1: measure a unit of capacity Run the job at a stable input rate and record throughput per worker: elements (or bytes) per second that one worker sustains, along with its CPU utilisation. This is the only honest basis for sizing, and it is job-specific — a pipeline doing a per-element external lookup has a completely different per-worker rate from one doing a local parse. ## Step 2: size for the spike, not the average With a per-worker rate, the arithmetic is direct. If the worst observed spike leaves a backlog of *B* elements and the SLA says it must be cleared within *T*, you need roughly `B / (T × per-worker rate)` workers *in addition* to the workers absorbing the ongoing input rate. Set `--max_num_workers` (Java: `--maxNumWorkers`) somewhat above that figure so the autoscaler has room, and treat it as a **guardrail against runaway cost**, not as a capacity target the job is expected to reach. Set `--num_workers` to a sensible starting size so the job does not spend the first minutes of every deploy scaling up from one. ## Step 3: remove the architectural constraint first Before tuning numbers, make sure scaling can actually happen quickly. Without Streaming Engine, a streaming job keeps its state on Persistent Disks attached to workers, provisioned at launch against the maximum worker count. Scaling then means redistributing state across disks — slow, disruptive, and paid for whether or not you reach the ceiling. Enabling Streaming Engine makes workers nearly stateless and the autoscaler far more responsive, both up and down. For a spiky workload this is usually a bigger win than any number you could tune, because the whole point of a spiky profile is that you want to *stop* paying at the quiet times. ## Step 4: check that the ceiling is reachable A maximum worker count is an upper bound, not a promise. Three things commonly stop a job from using it: - **Key-space parallelism.** Work in a stateful stage is partitioned by key. If your aggregation is keyed on ten tenant IDs, ten units of parallelism is what you get, no matter how many workers exist. The fix is upstream — a wider or salted key — not more machines. - **Downstream capacity.** More workers means more concurrent writes to the sink and more calls to any external service in a `DoFn`. If the sink throttles, adding workers converts a backlog problem into a retry storm. Rate-limit deliberately, batch writes, or size the sink to the peak worker count. - **Quota and machine availability.** CPU, IP address, and disk quotas in the region cap the fleet. Confirm the ceiling is grantable before relying on it. ## Step 5: know what the autoscaler is reacting to Dataflow's streaming autoscaler works from **backlog** and **CPU utilisation** signals — how much unprocessed input is queued and how hard the current workers are working. That is why backlog metrics deserve a place on the dashboard next to freshness and system latency, and why the autoscaling event log in the job UI is the first thing to read after a spike: it tells you whether the service tried to scale, how far it got, and whether it hit the ceiling. ## Step 6: consider the vertical axis too Horizontal autoscaling changes how *many* workers you have. If the symptom is a stage running out of memory rather than falling behind, more workers will not help; Dataflow Prime (`--dataflow_service_options=enable_prime`) adds vertical autoscaling and right-fitting so worker memory adapts to what the job actually uses. Recognising which axis a problem lives on is a large part of the judgment being tested. ## Step 7: decide when not to scale at all Sometimes the correct answer to a spiky backlog is to fix the pipeline rather than buy workers. A per-element synchronous call to an external API, an unbatched sink write, a hot key, or an unnecessary reshuffle all cap per-worker throughput; improving any of them lowers the worker count needed at every point on the curve. And some workloads should absorb the spike rather than chase it: if the SLA tolerates a ten-minute lag, letting a modest fleet burn the backlog down is cheaper than provisioning to clear it in sixty seconds. ## What to verify afterwards After the next real spike, check three things: did data freshness stay inside the SLA; did the autoscaling log show the job reaching, but not sitting at, the ceiling; and did the job actually scale back down afterwards. A job that never scales down is where streaming budgets die, and it usually points either to state pinned to disks or to a maximum set as a target rather than a guardrail.
- A Dataflow streaming job never reaches its configured maximum worker count despite a growing backlog. What do you check?First the key space: stateful work partitions by key, so a stage keyed on few values cannot use more parallelism regardless of fleet size. Then the autoscaling event log for what the service attempted, then quota and machine availability in the region, and then whether a throttled sink or external call is capping per-worker throughput rather than worker count.
- Why is Streaming Engine relevant to an autoscaling conversation at all?Without it, streaming state lives on Persistent Disks attached to workers and provisioned at launch against the maximum worker count, so scaling means moving state between disks — slow, disruptive, and paid for even at idle. Streaming Engine makes workers nearly stateless, which is what lets the autoscaler add and, crucially, remove workers quickly.
- When is the right answer to a spiky backlog not to scale out?When per-worker throughput is capped by the pipeline rather than by the fleet — a synchronous per-element external call, unbatched sink writes, a hot key, an avoidable reshuffle — or when the freshness SLA is loose enough that a modest fleet burning the backlog down costs less than provisioning to clear it in a minute.
- How do you tell a horizontal scaling problem from a vertical one?Falling behind with high CPU utilisation and a growing backlog across many workers is horizontal — you need more of them. Repeated out-of-memory failures or one memory-hungry stage is vertical, and the answer is worker sizing: Dataflow Prime adds vertical autoscaling and right-fitting so memory adapts to actual use rather than to a machine type you guessed.
saying these in an interview costs you the question
- Sets maximum workers very high and calls it capacity planning
- Ignores that key-space size caps usable parallelism
- Forgets the sink and external services must absorb peak worker count
- Treats the maximum worker count as a target rather than a guardrail
- Never checks whether the job scales back down after a spike