In an Apache Beam DoFn, what runs in setup versus start_bundle versus process?
answer
- it is nested, not one call per element
- expensive clients get built once, not per row
- the batch the runner retries has its own hooks
- one hook is only best effort
- setup, start_bundle, process, finish_bundle, teardown
basics
~20 sIn an Apache Beam DoFn, setup runs once per DoFn instance on the worker for expensive one-time init, start_bundle once before each bundle of elements, process once per element, finish_bundle after the bundle's elements, and teardown best-effort when the instance is discarded.
solid answer
~50 sA `DoFn` is the per-element function a `ParDo` applies, and its lifecycle is nested. `setup()` runs once when the runner instantiates the `DoFn` on a worker — that is where you open a client, a connection pool, or load a model. `start_bundle()` runs before each **bundle**, the runner-chosen batch of elements that is also the unit of retry; use it to reset per-bundle buffers. `process(element)` runs per element and yields output. `finish_bundle()` runs after the bundle's last element — flush buffered writes here. `teardown()` is **best effort**: the runner calls it when it discards the instance, but a crashed or preempted worker may never run it, so it cannot be your only cleanup path. One instance handles many bundles, and if a bundle fails the whole bundle is re-run, so anything you emit or write must tolerate re-execution.
code
python · 22 linesimport apache_beam as beam
class PostToApi(beam.DoFn):
def __init__(self, endpoint):
self.endpoint = endpoint # driver-side, must be picklable
def setup(self):
self.session = make_session() # once per instance, on the worker
def start_bundle(self):
self.buffer = [] # once per bundle
def process(self, element):
self.buffer.append(element) # once per element
yield element
def finish_bundle(self):
if self.buffer:
self.session.post(self.endpoint, json=self.buffer) # idempotent write
def teardown(self):
self.session.close() # best effort onlygo deeper
Recall that ParDo applies a DoFn to each element and that process() is the per-element method, and that beam.Map is a thin wrapper over ParDo.
Name all the hooks in order and say which scope each belongs to — instance, bundle, element — and give the reason expensive clients go in setup() rather than process() or the constructor.
Demonstrate the retry reasoning: bundles are re-run wholesale, so external side effects must be idempotent, in-memory counters lie, and teardown may never fire on a preempted worker.
Own the guidance your platform gives teams: where side effects are allowed in a pipeline at all, which sinks are safe to write from a DoFn, and when a dedicated IO connector should replace hand-rolled per-element writes.
## Where the lifecycle sits `ParDo` is Apache Beam's general-purpose element-wise transform; `beam.Map`, `beam.FlatMap` and `beam.Filter` are conveniences built on top of it. When you need per-element work with real setup, side outputs, or state, you write a `DoFn` subclass and pass it to `beam.ParDo`. The runner does not call your `DoFn` once per element in isolation. It creates instances on workers, groups elements into **bundles**, and drives a nested lifecycle around them. ## The five hooks, outermost first **`__init__`** runs in your *driver program*, on your machine, when you build the graph. Whatever you assign here is serialized (pickled, in Python) and shipped to workers. That is why you must not create a database client or an open socket in `__init__` — it will not serialize, or it will serialize into something dead. Constructor arguments should be plain configuration values. **`setup()`** runs once per `DoFn` instance, on the worker, before any element is processed. This is the correct home for expensive, reusable initialization: an API session, a connection pool, a loaded model file, a compiled regex table. Because one instance serves many bundles, this cost is amortized. **`start_bundle()`** runs once before the elements of a bundle. A bundle is a runner-decided group of elements — its size is not something you control or should assume. Use `start_bundle()` to reset per-bundle mutable state, typically an output buffer. **`process(element)`** runs once per element and returns or yields zero, one, or many outputs — that is why `ParDo` generalizes both map and filter. Extra arguments can be injected: `beam.DoFn.TimestampParam`, `beam.DoFn.WindowParam`, and side inputs passed by the caller. **`finish_bundle()`** runs after the last element of a bundle. Flush whatever `process()` buffered. A sharp edge: if you *emit output* from `finish_bundle()`, you must emit a `WindowedValue` with an explicit timestamp and window, because there is no current element to inherit them from. Most people only flush side effects here and emit nothing. **`teardown()`** is called when the runner discards the instance — close the connection you opened in `setup()`. It is explicitly **best effort**. A worker that is killed, preempted, or crashes will never call it. Never rely on `teardown()` for correctness, only for tidy resource release. ```python class PostToApi(beam.DoFn): def __init__(self, endpoint): self.endpoint = endpoint # plain config, gets serialized def setup(self): self.session = make_session() # once per instance, on the worker def start_bundle(self): self.buffer = [] def process(self, element): self.buffer.append(element) yield element def finish_bundle(self): if self.buffer: self.session.post(self.endpoint, json=self.buffer) def teardown(self): self.session.close() # best effort only ``` ## Bundles are the retry unit — and that is the real interview point If any element in a bundle fails, the runner re-runs **the whole bundle**, calling `start_bundle()`, `process()` on every element again, and `finish_bundle()` again. Consequences: - Side effects in `process()` or `finish_bundle()` can happen more than once. An external write must be idempotent — an upsert keyed by a stable id, or a write whose duplicate is harmless — or you will double-count on the first transient failure. - Counters you keep in Python variables are wrong under retries; use `beam.metrics.Metrics.counter` if you want observable counts, and treat even those as approximate. - State that survives across bundles in instance attributes is dangerous. Anything set in `setup()` should be read-only or a reusable resource; anything accumulating results belongs in `start_bundle()` so a retry starts clean. ## Common bugs this question is really probing 1. **Opening a connection per element.** Putting client construction in `process()` is the classic performance defect; it belongs in `setup()`. 2. **Opening a connection in `__init__`.** Unserializable, and a source of "works locally, fails on the service" errors. 3. **Mutating the incoming element.** Beam requires input elements to be treated as immutable; mutating one breaks retry correctness and the `DirectRunner` will flag it. 4. **Assuming bundle size.** You cannot ask for batches of exactly 500 by relying on bundles. If you need a controlled batch size, use `beam.BatchElements` or `GroupIntoBatches`, both of which are real Beam transforms designed for that. 5. **Relying on `teardown()` to flush data.** Flush in `finish_bundle()`; `teardown()` only releases resources. ## What to say "`__init__` on the driver and gets pickled, `setup()` once per instance on the worker for expensive clients, `start_bundle()` per bundle to reset buffers, `process()` per element, `finish_bundle()` to flush, `teardown()` best effort. The bundle is the retry unit, so every external side effect in there has to be idempotent."
- Why is a database connection created in the DoFn's constructor a problem?The constructor runs in the driver program, and the resulting object is serialized and shipped to workers. Live sockets and clients generally are not picklable, and even if they were, the deserialized copy would be useless on another machine. Store plain configuration in `__init__` and build the client in `setup()`, which runs on the worker.
- You need to send records to an API in batches of 500 — can you rely on bundle size?No. Bundle size is chosen by the runner and varies with the runner, the mode, and the moment; it is not tunable as a batch-size knob. Use `beam.BatchElements` or `GroupIntoBatches` to get controlled batch sizes, and still make the write idempotent because the enclosing bundle can be retried.
- If finish_bundle already flushed writes and the bundle then fails, what happens?The runner re-runs the entire bundle: `start_bundle`, every element through `process`, and `finish_bundle` again — so the flush happens a second time. That is why the external write must be idempotent, typically an upsert on a stable key or a deduplicated insert, rather than a blind append.
saying these in an interview costs you the question
- Says setup runs once per element
- Opens a database client in the DoFn constructor
- Treats teardown as a guaranteed flush point
- Assumes a bundle is a fixed, tunable batch size
- Assumes process runs exactly once per element