In Apache Beam, what is a PCollection and what properties does it guarantee?
answer
- the thing that flows between pipeline steps
- you never edit it in place
- its size may be finite or endless
- immutable distributed dataset with a coder
- bounded versus unbounded decides windowing
basics
~20 sA 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 sA `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 linesimport 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 fromgo deeper
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.
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.
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.
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