skip to content

Google Cloud Dataflow

Managed Apache Beam on GCP, running the same pipeline code in batch or streaming with autoscaling. Interviewers ask about windowing, watermarks and late data, because that is the part of the Beam model streaming will not let you avoid.

on this pageshow

explore

questions

12

In Apache Beam, what is a PCollection and what properties does it guarantee?

level: juniorimportance: must knowfreq 80%

answer

  1. the thing that flows between pipeline steps
  2. you never edit it in place
  3. its size may be finite or endless
  4. immutable distributed dataset with a coder
  5. bounded versus unbounded decides windowing

basics

~20 s

A PCollection is Apache Beam's immutable, distributed dataset that flows between transforms. It is unordered, either bounded or unbounded, carries a coder and a windowing strategy, and is never modified in place — every PTransform produces a new one.

solid answer

~50 s

A `PCollection` is the data abstraction in an Apache Beam pipeline: a potentially huge, distributed, **immutable** collection of elements that a `PTransform` consumes and produces. Three properties matter in interviews. First, immutability — applying a transform never changes the input; it returns a new `PCollection`, which is what lets a runner replay work after a failure. Second, boundedness — a `PCollection` is *bounded* (a finite file or table read) or *unbounded* (a continuous source such as a Pub/Sub subscription), and that flag drives whether windowing and triggers are required. Third, it has no useful ordering and no random access: you cannot index it, iterate it in your driver program, or rely on input order surviving a transform. Each `PCollection` also owns a coder (how elements are serialized between workers) and a windowing strategy, and belongs to exactly one `Pipeline`.

code

python · 10 lines
python
import apache_beam as beam

with beam.Pipeline() as p:
    lines = p | "Read" >> beam.io.ReadFromText("gs://bucket/in/*.txt")
    words = lines | "Split" >> beam.FlatMap(lambda line: line.split())
    counts = words | "Count" >> beam.combiners.Count.PerElement()
    counts | "Write" >> beam.io.WriteToText("gs://bucket/out/counts")

# lines, words and counts are separate PCollections;
# no step changed the collection it read from

go deeper

for a junior

Be ready to define PCollection, PTransform and Pipeline in one sentence each, and to say plainly that a transform returns a new PCollection instead of changing the old one.

for a middle

Explain why immutability is required rather than merely tidy: runners retry work in bundles, so mutating inputs makes a retry produce different results. Also explain the bounded/unbounded flag and what it forces you to do about windowing.

for a senior

Show you know the failure modes: element mutation that only breaks under retries, order-dependent logic, and missing or wrong coders that pass locally and fail on a distributed runner.

for a principal

Own the framing that Beam's value is one dataset abstraction covering batch and streaming, and be able to say where that unification actually pays off for a team versus where it costs you runner-specific tuning.

## The three nouns of the Beam model An Apache Beam program is built from three things. A `Pipeline` is the container object that holds the whole graph. A `PCollection` is a dataset flowing through that graph. A `PTransform` is an operation applied to one or more `PCollection`s to produce new ones. You build the graph in a *driver program* (Python, Java, Go), and the driver hands the finished graph to a **runner** — `DirectRunner` locally, `DataflowRunner` on Google Cloud, `FlinkRunner` on Flink — which actually executes it. The key mental shift for anyone coming from ordinary Python or Java collections: a `PCollection` is a *deferred handle* to data, not the data itself. When your driver runs, no elements exist yet. You are describing work. ```python import apache_beam as beam with beam.Pipeline() as p: lines = p | "Read" >> beam.io.ReadFromText("gs://bucket/in/*.txt") words = lines | "Split" >> beam.FlatMap(lambda l: l.split()) ``` `lines` and `words` are `PCollection`s. `len(words)` is meaningless; `words[0]` is meaningless; printing `words` in the driver prints an object description, not rows. ## Immutability A `PTransform` never edits its input. `beam.Map(f)` applied to `lines` leaves `lines` exactly as it was and yields a brand-new `PCollection`. You may apply many transforms to the same `PCollection` — that is how you branch a pipeline — and each branch is independent. This is not a stylistic preference. Runners rely on it: a distributed runner splits work into **bundles** and retries a failed bundle by re-running the transform on the same input elements. If a transform mutated its inputs, a retry would see different data and produce a different answer. That is why mutating an element inside a `DoFn` (for example appending to a list you received, or editing a dict in place and yielding it) is a genuine bug even when it appears to work locally — the `DirectRunner` actively checks for it, and on a distributed runner it produces intermittent wrong results. ## Bounded and unbounded Every `PCollection` is either bounded or unbounded, a property it inherits from the source that created it. - **Bounded**: a finite, known-size dataset — files in Cloud Storage, a BigQuery table, an in-memory `beam.Create([...])`. - **Unbounded**: a continuously arriving stream — a Pub/Sub subscription, a Kafka topic. It has no end, so a transform that must see "all" the data (an aggregation) cannot simply wait for completion. That is why unbounded data forces you to divide elements into **windows** and choose **triggers** that decide when a window emits a result. A bounded `PCollection` implicitly lives in one global window, so a `GroupByKey` on it just works. The same pipeline code can run in batch or streaming precisely because Beam models both as `PCollection`s and differs only in this flag and the windowing you attach. ## No order, no random access A `PCollection` has no meaningful element order. A runner distributes elements across workers and processes them in whatever order is convenient; the `DirectRunner` deliberately shuffles order so that order-dependent code fails on your laptop rather than in production. If order matters, it must be expressed as data — a timestamp field, a sequence number — and enforced by a transform, never assumed. Every element also carries an **event timestamp** and a window assignment. For a batch file read the timestamp defaults to a minimum value; for a streaming source it usually comes from the message (for example a Pub/Sub publish time or a field you extract). Timestamps are what windowing reasons about. ## Coders Because elements must travel between workers and be persisted during a shuffle, each `PCollection` has a **coder**: the rule for turning an element into bytes and back. Beam infers coders for common types and for registered schemas; for a custom class you may need to register one. A missing or wrong coder is a classic "works on `DirectRunner`, fails on the service" bug, along with a `DoFn` that closes over something unpicklable. ## What to say in an interview "A `PCollection` is Beam's immutable distributed dataset. Transforms produce new ones rather than editing in place, which is what makes bundle retries safe. It is bounded or unbounded — that decides whether I need windowing — it carries per-element timestamps and a coder, and it has no order and no random access, so I never write code that depends on either."

  • Why can't you just call len() or index into a PCollection in your driver program?
    Because the driver only builds a graph — no elements exist when the driver runs, and at execution time the data is spread across many workers with no global ordering. Anything you need from the whole collection must be produced by a transform, for example `beam.combiners.Count.Globally()`, and then consumed downstream or fed back as a side input.
  • What does an element's event timestamp mean, and where does it come from?
    It is the logical time the event happened, and windowing and watermarks reason about it rather than about wall-clock processing time. Streaming sources usually supply it (a Pub/Sub publish time, or a field you extract and set with `beam.window.TimestampedValue`). Bounded file reads assign a default minimum timestamp, which is fine because they sit in one global window.
  • Can the same pipeline code run over both a bounded and an unbounded PCollection?
    Largely yes — that is the point of the unified model. The transform code is identical; what changes is the source and the windowing strategy. A global-window aggregation that terminates in batch would never emit on an unbounded input, so streaming versions need explicit windows and triggers, and some IO sinks only support one mode.

saying these in an interview costs you the question

  • Calls a PCollection an in-memory list you can index
  • Mutates elements inside a DoFn and yields them
  • Assumes input order survives a transform
  • Says PCollections exist only in batch pipelines
  • Thinks a transform edits its input collection in place

context

open as a page

In Apache Beam, why does CombinePerKey handle a hot key better than GroupByKey?

level: middleimportance: must knowfreq 68%

basics

~20 s

GroupByKey sends every value for a key across the shuffle and materializes them all on one worker, so a hot key becomes a straggler or an out-of-memory failure. CombinePerKey pre-aggregates with an associative, commutative CombineFn, so only small accumulators are shuffled and merged.

open as a page

In Dataflow, what is the difference between draining and cancelling a streaming job?

level: middleimportance: must knowfreq 70%

basics

~20 s

Draining a Dataflow streaming job stops it from reading new input, advances the watermark to infinity so every open window fires and in-flight data is written to sinks, then finishes. Cancelling stops immediately and discards buffered state and in-flight data.

open as a page

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

level: seniorimportance: must knowfreq 60%

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.

open as a page

In an Apache Beam DoFn, what runs in setup versus start_bundle versus process?

level: middleimportance: should knowfreq 58%

basics

~20 s

In 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.

open as a page

In Dataflow, what does enabling Streaming Engine change about how a streaming job runs?

level: middleimportance: should knowfreq 55%

basics

~20 s

Streaming Engine moves pipeline state, timers and shuffle off the worker VMs into the Dataflow backend service. Workers become nearly stateless, so they need much smaller disks and less CPU and memory, and horizontal autoscaling can react far faster.

open as a page

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

level: seniorimportance: should knowfreq 52%

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.

open as a page

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

level: seniorimportance: should knowfreq 44%

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.

open as a page

In the Dataflow job UI, what do the data freshness and system latency graphs measure?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Data freshness is the gap between now and the job's output watermark — how far behind real time fully-processed data is. System latency is the longest time any element currently in the pipeline has spent being processed or waiting to be processed.

open as a page

How would you size horizontal autoscaling for a Dataflow streaming job with a spiky Pub/Sub backlog?

level: principalimportance: should knowfreq 35%

basics

~20 s

Measure 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.

open as a page

In a Dataflow pipeline, what does PubsubIO's withIdAttribute change about deduplication?

level: middleimportance: nice to knowfreq 30%

basics

~20 s

By default the Dataflow runner deduplicates Pub/Sub messages using the service-assigned message ID. withIdAttribute tells PubsubIO to dedupe on a publisher-set message attribute instead, so two separate publishes of the same logical event collapse into one.

open as a page

For a Dataflow job other teams launch, how do classic and Flex templates differ?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

A Dataflow classic template stages a pre-built job graph in Cloud Storage, so per-launch parameters must be read through ValueProvider and the graph shape is fixed. A Flex template packages the pipeline in a container image and builds the graph at launch, with real parameter values.

open as a page