Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Apache Beam model

Batch & Streaming · Streaming Design

Apache Beam model

Mediumbatch-streaming-39
apache-beamdataflowwindowingstream-processing

Question

What is the Apache Beam programming model?

Solution

The Apache Beam programming model offers a unified API for building both batch and streaming pipelines through core abstractions called PCollections and PTransforms. It organizes distributed computation around answering four questions: what is computed, where in event time it is computed, when in processing time results are emitted, and how multiple results relate. Developers write pipeline logic once and execute it across distinct distributed engines using runners like Google Cloud Dataflow, Apache Flink, and Apache Spark.

Core abstractions: PCollections and PTransforms

At the center of Beam are two fundamental primitives:

  • PCollection: Represents a distributed, potentially unbounded dataset across cluster nodes. In batch mode, a PCollection is bounded by fixed input files; in streaming mode, it is unbounded and continuously updated.
  • PTransform: Represents a processing operation applied to one or more PCollections, such as ParDo for record-by-record transformations, GroupByKey for shuffling data, and Combine for distributed aggregations.

The four questions framework

Beam clarifies stream and batch semantics by separating processing requirements into four explicit decisions:

  • What is being computed: Defined by user transformations, such as computing revenue sums or filtering invalid records.
  • Where in event time: Defined by windowing strategies, grouping unbounded data streams into fixed, sliding, or session time windows.
  • When in processing time: Defined by triggers, dictating when an operator emits window calculation results relative to watermark progression or processing time elapsed.
  • How results relate: Defined by accumulation mode, deciding whether successive pane firings overwrite prior results with accumulated totals or emit only incremental deltas.

Runner portability in production

Beam decouples developer business logic from backend infrastructure. You can develop and test pipelines locally using the DirectRunner, run massive nightly backfills on an Apache Spark cluster, and deploy continuous streaming workloads onto Google Cloud Dataflow or Apache Flink.

Writing once and selecting your runner gives teams architectural flexibility to migrate workloads across cloud providers or execution engines without rewriting underlying pipeline code.

PreviousNext