skip to content

In Apache Beam, what actually changes when you swap DirectRunner for DataflowRunner?

level: seniorimportance: should knowfreq 52%

answer

  1. the graph is the same either way
  2. one runner is a correctness harness, not a small cluster
  3. laptop-to-service breakage is about packaging
  4. serialization, dependencies, file paths
  5. fusion and autoscaling only exist on one side

basics

~20 s

The pipeline code and graph stay identical; only execution changes. DirectRunner runs everything in one local process and deliberately stresses model rules, while DataflowRunner submits the graph to the managed service, which provisions workers, fuses steps and needs project, region and staging locations.

solid answer

~50 s

Apache Beam separates the SDK from execution: the same `Pipeline` of `PCollection`s and `PTransform`s is submitted to whichever runner you name, so **your transform code does not change at all**. What changes is everything around it. `DirectRunner` executes locally in one process for testing, and deliberately makes life hard — it processes elements in arbitrary order, checks that you do not mutate input elements, and round-trips elements through their coders — so model violations fail on your laptop. `DataflowRunner` serializes the graph and hands it to the Google Cloud Dataflow service, which provisions workers, fuses adjacent steps into stages, autoscales, and runs your SDK code in worker containers; it requires `--project`, `--region`, and staging and temp Cloud Storage locations. The bugs that only appear after the swap are unpicklable `DoFn` state, missing worker dependencies, unregistered coders, and reliance on the local filesystem.

code

bash · 11 lines
bash
# local correctness run
python pipeline.py --runner=DirectRunner

# same pipeline code, submitted to the managed service
python pipeline.py \
  --runner=DataflowRunner \
  --project=my-project \
  --region=us-central1 \
  --temp_location=gs://my-bucket/tmp \
  --staging_location=gs://my-bucket/staging \
  --requirements_file=requirements.txt

go deeper

for a junior

Know that the runner is chosen with a flag, that DirectRunner runs locally for testing, and that DataflowRunner submits the job to Google Cloud.

for a middle

Explain that the graph is unchanged and describe what DirectRunner deliberately does — arbitrary ordering, immutability and coder checks — and what the service adds: provisioning, fusion, shuffle and autoscaling.

for a senior

Name and diagnose the laptop-to-service failure classes: unpicklable DoFn state, missing worker dependencies, local filesystem paths, missing coders, and stale global state, plus the required project, region and staging options.

for a principal

Own the portability position honestly: what runner independence is worth to your organization, where the capability matrix and operational differences make a migration real work, and whether that optionality justifies the abstraction cost.

## Portability is the premise Apache Beam's design splits the **SDK** (how you express a pipeline) from the **runner** (what executes it). Your driver program builds a graph of `PTransform`s over `PCollection`s; the runner receives that graph. `DirectRunner`, `DataflowRunner`, `FlinkRunner`, `SparkRunner` and others all consume the same graph. In Python you usually pick one on the command line: ```bash python pipeline.py --runner=DirectRunner python pipeline.py \ --runner=DataflowRunner \ --project=my-project \ --region=us-central1 \ --temp_location=gs://my-bucket/tmp \ --requirements_file=requirements.txt ``` So the honest first half of the answer is: **nothing in the pipeline changes**. No `ParDo` is rewritten, no `CombinePerKey` behaves differently, no windowing is redefined. That is the portability claim, and interviewers ask this to see whether you can separate the SDK from the service. ## What DirectRunner actually is `DirectRunner` is not "a small Dataflow". It is a *model-conformance* runner. It runs the pipeline locally, typically in one process, over data small enough for a test, and it deliberately behaves adversarially so that assumptions you should not be making break immediately: - It processes elements in an **arbitrary order**, so order-dependent logic fails locally rather than intermittently in production. - It **checks immutability**, catching a `DoFn` that mutates an element it received or one it already emitted. - It **encodes and decodes elements through their coders**, so a type without a usable coder fails at once. - It serializes user functions, catching some unpicklable closures. Because of these checks, `DirectRunner` is slow and is for correctness testing, not for benchmarking. Do not draw performance conclusions from it. ## What DataflowRunner adds Submitting with `DataflowRunner` turns your driver into a *job submitter*. The driver builds the graph, uploads dependencies and the pipeline package to your staging location, and creates a Dataflow job. From there the service owns execution: - **Worker provisioning** — VMs are created in your project and region, running the Beam SDK harness in containers. - **Fusion and optimization** — adjacent element-wise steps are fused into a single stage so elements are not materialized between them. This is why the step names in the monitoring UI often do not map one-to-one to your transforms, and why an artificial fusion break (`beam.Reshuffle()`) is sometimes used to force parallelism. - **Shuffle and grouping as a service.** - **Autoscaling and work rebalancing** across workers. It also introduces a whole class of concerns your laptop never had: IAM permissions for the worker service account, network and subnet configuration, staging and temp buckets, and dependencies that must be present *on the workers*. ## The bugs that appear only after the swap This is the part interviewers actually want. 1. **Unpicklable `DoFn` state.** Anything assigned in `__init__` is serialized to workers. A database client, a file handle, a thread, or a closure over a module-level object may fail to serialize. Fix: keep configuration in `__init__` and build resources in `setup()`. 2. **Missing dependencies on workers.** The import that resolves locally is not installed in the worker container. Fix: `--requirements_file`, `--setup_file` for a local package, or a custom container image. 3. **Local filesystem assumptions.** Reading `./config.json` or writing `/tmp/out.csv` works locally and silently writes to a random worker on the service. Fix: use Cloud Storage paths and Beam's filesystems layer. 4. **Missing or wrong coders** for custom types, which `DirectRunner` may catch but a lax local path can hide. 5. **Global mutable state** shared between elements, which happens to be consistent in one local process and is not across workers. 6. **Order and singleton assumptions** — code that worked because the local run happened to produce one bundle. ## Portability is not uniformity The same code compiles against every runner, but not every runner implements every corner of the model to the same degree — Beam publishes a capability matrix precisely because support for things like certain state and timer features or specific IO connectors varies. And even where semantics match, operations do not: scaling behaviour, monitoring, and failure handling are runner-specific. Treat portability as "my transform logic is not locked in", not as "I can move a production pipeline between runners with no work". ## What to say "The graph is identical — that is the point of Beam. `DirectRunner` is a local correctness harness that shuffles order, checks immutability and round-trips coders so model violations fail on my laptop. `DataflowRunner` ships the graph to the service, which provisions workers, fuses steps, shuffles and autoscales, and needs project, region and GCS staging. What breaks in the move is serialization of `DoFn` state, worker dependencies, local filesystem paths, and coders — not the pipeline logic."

  • A pipeline passes locally but fails on Dataflow workers with an ImportError — where do you look?
    Worker packaging. The driver machine has the library, the worker container does not. Supply it with `--requirements_file` for pip dependencies, `--setup_file` for a local package that needs building, or a custom SDK container image. Confirm the module is a real dependency and not something you only had installed ad hoc in your dev environment.
  • Why does the Dataflow monitoring UI show fewer steps than your pipeline has transforms?
    Because the service fuses adjacent element-wise steps into a single stage so elements never have to be materialized between them. It is an optimization, but it also means the fused stage inherits the parallelism of its input; when that limits throughput, inserting `beam.Reshuffle()` creates a deliberate fusion break so downstream work can spread across more workers.
  • Does the same Beam pipeline behave identically on FlinkRunner and DataflowRunner?
    The transform semantics are meant to match, but support is not uniform — Beam publishes a capability matrix because runners differ in coverage of features such as certain state and timer facilities and available IO connectors. Operations differ far more: scaling, monitoring, checkpointing and failure handling are runner-specific. Portability means the logic is not locked in, not that a migration is free.

saying these in an interview costs you the question

  • Treats DirectRunner results as a performance benchmark
  • Says pipeline code must be rewritten per runner
  • Opens clients in the DoFn constructor and blames the service
  • Assumes local file paths work on workers
  • Claims full feature parity across all Beam runners

context