skip to content

In Apache Beam, when does a large side input become the bottleneck in a ParDo?

level: seniorimportance: should knowfreq 44%

answer

  1. it is the broadcast side of a join
  2. every worker needs all of it, not a slice
  3. scaling up multiplies the cost
  4. AsList and AsDict live in worker memory
  5. when both sides are big, shuffle instead

basics

~20 s

A side input is broadcast in full to every DoFn invocation, so a large one must be materialized, distributed to every worker, and held in worker memory — with AsList or AsDict fully in memory. Past that point a CoGroupByKey join is the right shape.

solid answer

~40 s

A **side input** in Apache Beam is an extra `PCollection` made available whole to every invocation of a `DoFn`, passed with `beam.pvalue.AsDict`, `AsList`, `AsSingleton` or `AsIter`. It is the broadcast, map-side join: perfect for a small lookup table such as currency rates or a country dimension. It stops being cheap when the collection is large, because the runner must materialize it, ship it to every worker, and keep it accessible there — and `AsList`/`AsDict` build the whole structure in worker memory, so a big one causes memory pressure or repeated fetching. It is also a **barrier**: the consuming `ParDo` cannot process elements for a window until the side input for that window is fully computed. Once both sides are genuinely large, switch to a shuffle-based join with `CoGroupByKey`.

code

python · 10 lines
python
import apache_beam as beam

# broadcast join: small reference table as a side input
rates = (p | "ReadRates" >> beam.io.ReadFromText("gs://ref/fx.csv")
           | "Parse" >> beam.Map(parse_kv))          # -> (ccy, rate)

converted = (orders
             | "Convert" >> beam.Map(
                 lambda order, fx: {**order, "usd": order["amount"] * fx[order["ccy"]]},
                 fx=beam.pvalue.AsDict(rates)))

go deeper

for a junior

Know that a side input is extra data passed to a ParDo alongside each element, and that beam.pvalue.AsDict or AsList is how you pass it.

for a middle

Explain that the whole collection is visible per invocation, name the four view wrappers and what each materializes, and say why that makes side inputs a broadcast join for small reference data only.

for a senior

Diagnose the real symptoms — OOM in a trivial step, a stall before the step starts, scaling up making things worse, periodic streaming latency — and choose between shrinking the side input, changing the view, or moving to CoGroupByKey.

for a principal

Own the enrichment strategy across pipelines: which reference datasets are allowed to be broadcast, when a shared key-value store beats embedding lookups in every job, and how refresh cadence for reference data is agreed and monitored.

## What a side input is Every `PCollection` a `ParDo` reads element-by-element is its *main input*. A **side input** is an additional `PCollection` that the transform can consult *in its entirety* on every element. In Python you wrap it at the call site: ```python rates = (p | "ReadRates" >> beam.io.ReadFromText("gs://ref/fx.csv") | "Parse" >> beam.Map(parse_kv)) converted = (orders | "Convert" >> beam.Map( lambda order, fx: {**order, "usd": order["amount"] * fx[order["ccy"]]}, fx=beam.pvalue.AsDict(rates))) ``` The wrappers describe the view your function receives: - `AsSingleton` — the collection must contain exactly one element; you get that value. - `AsIter` — an iterable over the elements. - `AsList` — a real list, fully materialized. - `AsDict` — a dict built from `(key, value)` pairs; keys must be unique. Conceptually this is a **broadcast join**: the small side is replicated everywhere so the big side never has to be shuffled. ## Why size is the whole question Three costs scale with the side input. **Materialization.** The side input is itself the output of a subgraph. Before the consuming `ParDo` can run, that subgraph must complete and its result be written somewhere the workers can read. **Distribution.** Every worker that runs the consuming `DoFn` needs access to the entire side input, not a slice of it. Total work is roughly side-input size times number of workers, which is exactly the term that grows when a job autoscales up. **Worker memory.** `AsList` and `AsDict` construct the whole structure in the worker process. A side input that comfortably fits on your laptop under `DirectRunner` can be the reason a distributed job dies with an out-of-memory error, and the failure mode is confusing because the *main* input looks innocent. `AsIter` avoids building one giant Python object but still requires access to the full collection. There is no published size at which this flips, and you should not quote one in an interview. The honest answer is a shape rule: side inputs are for reference data small enough to sit comfortably in worker memory alongside your other state; once the lookup side approaches the scale of the main input, use a shuffle-based join. ## The barrier property A side input is also a synchronization point. The consuming `ParDo` cannot emit results for a window until the side input for that window is ready. In batch this means the side-input branch gates the main branch — an expensive side-input computation delays everything downstream of it. In streaming, side inputs are **per window**: as the main input advances into a new window, the runner needs the side-input view for that window, and main-input elements wait until it is available. A side input derived from a slowly updating reference stream is a common source of "my streaming pipeline stalls periodically" reports. ## The alternative: CoGroupByKey When both sides are large, the correct primitive is a shuffle join: ```python joined = ({"orders": keyed_orders, "customers": keyed_customers} | beam.CoGroupByKey()) ``` `CoGroupByKey` groups several keyed `PCollection`s by the same key and hands you, per key, the values from each tagged collection. Both sides are shuffled — more network than a broadcast — but nothing has to fit on a single worker, and the cost scales with data rather than with data times worker count. The trade is the classic one: broadcast join when one side is small, shuffle join when neither is. ## Diagnosing it in practice Symptoms that point at a side input: - Workers OOM in a step that does trivial per-element work. - The job stalls before the consuming step even starts, while an unrelated-looking branch is still running. - Scaling the job up makes it *worse*, because every new worker must pull the whole side input again. - A streaming pipeline's latency spikes at regular intervals matching the side-input window. Fixes, in order of preference: shrink the side input (project only the columns you need, filter to the keys that actually occur, pre-aggregate it); switch `AsList` to `AsIter` or `AsDict` keyed lookups so you are not holding a list you scan linearly; move to `CoGroupByKey`; or, for genuinely huge reference data with point lookups, query an external key-value store from `setup()` in the `DoFn` and accept the per-lookup latency. ## What to say "A side input is Beam's broadcast join — the whole collection is visible to every `DoFn` call, wrapped with `AsDict`, `AsList`, `AsSingleton` or `AsIter`. It is great for small reference data and bad when it is large, because it is materialized, copied to every worker, and with `AsList`/`AsDict` held in worker memory, and because it gates the consuming step until it is ready per window. Once both sides are big I use `CoGroupByKey` instead."

  • How much of a side input is visible to a single DoFn invocation?
    All of it, for the relevant window. That is the defining property — unlike a main input element, which is one record, the side input view is the whole collection, which is why every worker must have access to the entire thing and why size is the deciding factor.
  • How do side inputs behave in a streaming pipeline?
    They are per window. The runner provides the side-input view corresponding to the main element's window, and main-input elements cannot be emitted until that view is ready. A slowly refreshing reference stream therefore shows up as periodic latency spikes on the consuming step, timed to the side input's windowing.
  • You need to enrich a huge event stream from a huge customer table — what shape do you use?
    A shuffle join with `CoGroupByKey`: key both collections on customer id and group them together, then emit the joined rows. Both sides pay shuffle cost, but neither has to fit on one worker and cost scales with data rather than with data times worker count. Reserve side inputs for reference data that stays small.

saying these in an interview costs you the question

  • Says a side input is only the elements matching the key
  • Treats AsList as free regardless of collection size
  • Adds workers to fix a side-input memory failure
  • Ignores that side inputs gate the consuming step
  • Uses a side input to join two large collections

context