Apache Beam is a programming model and a set of SDKs for writing data-processing pipelines once and running them on several execution engines. The same pipeline can read a bounded file or an unbounded stream, and the same code can run on Google Cloud Dataflow, Apache Flink, Apache Spark or a local runner. Beam does not execute anything itself. It defines what a pipeline means, and a runner translates that meaning into work on its engine.
This article is about the architecture that makes that possible and the consequences it has for your code: how a pipeline becomes a language-neutral graph, how runners talk to your code through the Fn API, why your DoFn sees elements in bundles, where shuffles happen, how per-key state and timers work, and how Java connectors can be used from Python. Event-time windows, watermarks and triggers are the core of Beam's streaming semantics, but they are covered in windowing in stream processing; here they appear only where they affect the architecture.
The model in five abstractions
A Pipeline is the whole computation, a directed acyclic graph. A PCollection is an immutable, possibly unbounded collection of elements; every element carries a value, an event timestamp, the window or windows it belongs to, and pane information that says which firing of a window produced it. A PTransform is an operation from PCollections to PCollections, and composite transforms are made of other transforms, so a whole IO connector is one node you can expand into dozens.
Everything is built from a few primitives. ParDo runs a user function, a DoFn, on each element and can emit zero or more outputs, including to multiple tagged outputs. GroupByKey collects all values with the same key and window. Flatten merges PCollections, and Impulse and splittable DoFns create data from nothing. Combine, CoGroupByKey, Reshuffle and every IO are composites over these. Side inputs let a DoFn read a whole PCollection, windowed to match the main input, as a map or list.
The key design choice is that bounded and unbounded data use the same abstractions. A batch job is a stream that ends: its watermark jumps to infinity when the input is exhausted, and every window fires once. That is why a pipeline tested on a file usually works on a stream, and also why you must still think about lateness when you switch.
The portability architecture
Beam has SDKs in Java, Python and Go, and runners written mostly in Java. The portability framework connects them with two protocol layers defined in protocol buffers and gRPC.
The Runner API is the pipeline representation: transforms, PCollections, coders, windowing strategies and environments, each identified by a URN. When you call pipeline.run(), the SDK serialises the graph into this proto and hands it to a job service. The Fn API is how a runner executes user code it cannot run itself. The runner starts SDK harness processes, typically containers for the pipeline's environment, and talks to them over separate channels: control (process this bundle, split it, finalize it), data (elements in and out), state (read and write user state and side inputs), and logging. The runner owns shuffling, state storage, timers, watermarks and checkpoints; the harness owns your DoFns and coders.
This split explains several practical behaviours. A Python pipeline on Dataflow or Flink runs your Python code in a harness next to a Java runner, so every element crossing between them is encoded with its coder. Elements pass between fused steps inside a harness without encoding. And a pipeline's dependencies are those of the harness environment, so the container image or requirements file you ship is part of the pipeline's correctness, not packaging trivia.
A worked pipeline
Here is a streaming pipeline in Python that counts purchases per product per minute from Pub/Sub, using the event time carried in a message attribute. It is short, but it exercises every architectural layer: an unbounded source, ParDo, windowing, a combine that becomes a shuffle, and a sink.
import json
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.transforms import window, trigger
opts = PipelineOptions(streaming=True)
with beam.Pipeline(options=opts) as p:
(p
| "Read" >> beam.io.ReadFromPubSub(
subscription="projects/shop/subscriptions/purchases",
timestamp_attribute="event_ts") # event time, not arrival time
| "Parse" >> beam.Map(json.loads)
| "Key" >> beam.Map(lambda e: (e["product_id"], 1))
| "Window" >> beam.WindowInto(
window.FixedWindows(60),
trigger=trigger.AfterWatermark(late=trigger.AfterCount(1)),
accumulation_mode=trigger.AccumulationMode.ACCUMULATING,
allowed_lateness=600)
| "Count" >> beam.CombinePerKey(sum)
| "Format" >> beam.MapTuple(lambda k, n: {"product_id": k, "purchases": n})
| "Write" >> beam.io.WriteToBigQuery(
"shop:analytics.purchases_per_minute",
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND))When it runs, the runner fuses Read, Parse, Key and the pre-combine half of Count into one stage. Combine is lifted: each worker computes partial sums per key before the shuffle, so only one partial count per key per window per bundle crosses the network instead of every purchase. The GroupByKey inside CombinePerKey is the shuffle boundary; after it, the merge half of Count, Format and Write form a second fused stage. Because the window accumulates and allows late firings, BigQuery receives an on-time row per product per minute and then an updated, larger row for each late element. Downstream queries must take the latest pane, or they double-count.
Fusion, bundles and why side effects repeat
Runners optimise the graph before running it. Fusion merges consecutive element-wise steps into one stage so elements flow through function calls instead of being written out between steps. Fusion stops at GroupByKey and at other points where data must be redistributed. It has a cost: if a cheap step fans out one element into a million, the fused stage keeps all of that work on the worker that read the original element. Insert beam.Reshuffle() after a large fan-out to break fusion and spread the work.
Within a stage, the runner hands elements to the harness in bundles. A bundle is the unit of processing and of failure: if anything in it fails, the runner retries the whole bundle, possibly on another worker, and discards its outputs and state changes. Bundle size is the runner's choice, and you should not depend on it. DoFn.setup runs once per DoFn instance and is where clients and connection pools belong; start_bundle and finish_bundle bracket a bundle and are where you buffer and flush batched writes.
The consequence for correctness is the most important thing to understand about Beam. The runner can make its own state and outputs effectively exactly-once, because it commits a bundle's results atomically. It cannot undo a side effect your DoFn performed, such as an HTTP call or a database insert. Those can happen more than once when a bundle is retried. Make external writes idempotent with deterministic keys, or use a sink that implements idempotent commits. Exactly-once processing covers the general pattern.
Per-key state and timers
A stateful DoFn gets durable state scoped to each key and window, plus timers that call back into the DoFn at an event time or processing time. The runner stores the state, checkpoints it, and routes all elements of a key to the same place, so the DoFn can implement logic that windowed combines cannot express: deduplication, sessionisation with custom rules, rate limiting or buffering until a condition holds.
from apache_beam.coders import BooleanCoder
from apache_beam.transforms.timeutil import TimeDomain
from apache_beam.transforms.userstate import ReadModifyWriteStateSpec, TimerSpec, on_timer
class DedupeFn(beam.DoFn):
# Emit the first event per key; forget the key one hour of event time later.
SEEN = ReadModifyWriteStateSpec("seen", BooleanCoder())
EXPIRY = TimerSpec("expiry", TimeDomain.WATERMARK)
def process(self, element,
ts=beam.DoFn.TimestampParam,
seen=beam.DoFn.StateParam(SEEN),
expiry=beam.DoFn.TimerParam(EXPIRY)):
key, event = element
if seen.read():
return # duplicate: drop it
seen.write(True)
expiry.set(ts + 3600) # event-time timer
yield event
@on_timer(EXPIRY)
def expire(self, seen=beam.DoFn.StateParam(SEEN)):
seen.clear() # bounded state per key
# events | beam.Map(lambda e: (e["event_id"], e)) | beam.ParDo(DedupeFn())Two rules keep stateful DoFns healthy. Always clear state with a timer or window expiry, or state grows forever with the key space. And remember that state is per key, so a hot key serialises on one worker; state cannot be parallelised within a key. Joins built on state are covered in stateful stream joins.
Splittable DoFns and IO
Sources are not a special primitive in modern Beam. A splittable DoFn processes one element, such as a file name or a Kafka partition, over a restriction, such as a byte range or offset range. The runner can ask the harness to split the remaining restriction while it is being processed, and hand the remainder to another worker. That gives two capabilities: dynamic work rebalancing in batch, where a straggling file range is split so idle workers can take part of it, and checkpointing in streaming, where an unbounded read stops at a resumable position and reports a watermark estimate. Most users never write one, but knowing they exist explains why a slow file read can be rebalanced while a slow element inside an ordinary ParDo cannot.
Cross-language transforms
Because pipelines are protos with URNs, one SDK can use transforms implemented in another. When a Python pipeline uses apache_beam.io.kafka.ReadFromKafka, the Python SDK calls an expansion service, a Java process, which expands the transform into its Java sub-graph and returns it with a Java environment. The runner then starts both a Java and a Python harness and passes elements between them using portable coders. The cost is an extra environment to build and version, and the constraint that elements crossing the language boundary must use coders both sides understand, such as rows with a schema, not arbitrary pickled Python objects.
Runners differ more than the model suggests
| Runner | Typical use | Points to verify |
|---|---|---|
| Dataflow | Managed batch and streaming on Google Cloud | Autoscaling, update and drain semantics, service-side shuffle and state |
| Flink | Self-managed streaming on your own clusters | Checkpoint interval, state backend, savepoint-based upgrades |
| Spark | Batch on existing Spark clusters | Streaming feature coverage, which is narrower than Flink's |
| Direct / local | Unit tests and development | Checks model rules, not production performance |
The Beam capability matrix on the project site lists which features each runner supports. Check it for anything beyond the basics, such as timers in merging windows, splittable DoFn checkpointing or a particular trigger, before you commit to a runner. If you choose Flink, its own checkpointing governs recovery, as described in the Flink article; on Google Cloud, see Dataflow.
Failure modes
- Poison element. One record that throws makes its bundle retry indefinitely in streaming, stalling the stage and its watermark. Catch exceptions in the DoFn and route bad records to a dead-letter output with a tagged output.
- Hot key. GroupByKey and state are per key, so one dominant key bottlenecks a single worker. Salt the key, combine partially, then merge.
- Non-deterministic coders. GroupByKey compares encoded keys, so it requires a deterministic key coder. Java rejects a non-deterministic key coder when the pipeline is built; Python can fail at run time on a key type such as a dict. Use tuples, strings or schema rows as keys.
- Stuck watermark. An idle or lagging input partition holds the watermark back, so windows never fire. Check per-source watermark metrics before blaming the trigger.
- Duplicate side effects. External calls made in a DoFn repeat on retry. Make them idempotent.
- Incompatible update. Renaming steps or changing coders breaks in-place update of a streaming job. Give transforms stable names and plan for drain and replace.
Trade-offs: when Beam is the right layer
Beam pays off when you want one programming model for batch and streaming, when you run on Dataflow, or when portability between engines matters. Its costs are abstraction and lag: engine-specific features appear later or not at all, debugging crosses the SDK harness boundary, and performance tuning still requires knowing the runner. A team that runs only Flink and wants its lowest-level APIs is often better served by Flink directly. A team that wants a managed service, a Python-first API and the same code for backfills and live processing is Beam's core audience.
What to do next
- Write a small batch pipeline with the local runner and unit-test each transform with
TestPipelineandassert_that. - Switch it to streaming with event timestamps from the source and read windowing to choose windows, triggers and allowed lateness.
- Mark every external write in a DoFn and make it idempotent, and add a dead-letter output for parse failures.
- Look at the job graph in your runner's UI, find the fused stages and shuffle boundaries, and add Reshuffle after any large fan-out.
- Check the capability matrix for every feature you use on your chosen runner.
- Give every transform a stable name and rehearse an update or drain-and-replace of the streaming job before going live.