Cloud Dataflow: Real-Time Stream Processing with Apache Beam

Cloud Dataflow is GCP's managed Apache Beam service for batch and streaming data pipelines. This guide covers the Beam programming model, windowing and triggers, stateful processing, Dataflow Flex Templates, and production pipeline operations with metrics and alerting.

Streaming data pipelines are where data engineering gets genuinely hard. With batch jobs, you process a fixed dataset, verify the results, and retry if something goes wrong. With streaming pipelines, data arrives continuously, events arrive late or out of order, and "retry if something goes wrong" means deciding exactly how to handle partially processed state.

Cloud Dataflow is GCP's managed service for Apache Beam — a unified programming model for both batch and streaming pipelines. "Unified" means the same pipeline code can process a historical BigQuery backfill or a live Pub/Sub stream by changing the runner configuration, not the pipeline logic.

This guide covers Apache Beam's core abstractions (PCollections, transforms, windows, triggers), stateful processing for advanced streaming patterns, and operational best practices for running Dataflow pipelines in production.

Apache Beam Core Concepts

PCollection: A distributed dataset. In batch mode, it's a finite dataset. In streaming mode, it's an unbounded sequence of elements arriving over time. Transforms take PCollections as input and produce PCollections as output.

PTransform: A data processing operation. Beam includes built-in transforms (ParDo, GroupByKey, Combine, Flatten, Partition) and you can compose them into higher-level transforms.

Pipeline: A directed acyclic graph of PCollections and PTransforms. The Pipeline Runner (Cloud Dataflow, DirectRunner, etc.) executes the graph.

Window: In streaming, grouping operations need a time boundary. Windows define how elements are grouped by time: fixed windows, sliding windows, or session windows.

Trigger: Specifies when to emit results for a window. The default trigger waits for the watermark to pass the end of the window, then emits once. You can trigger early (for speculative results) or late (for handling late data).

Your First Dataflow Pipeline

# pipeline.py
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions, GoogleCloudOptions

def parse_event(element: str) -> dict:
    """Parse a JSON event from Pub/Sub."""
    import json
    try:
        event = json.loads(element)
        return event
    except Exception:
        return None

def enrich_event(event: dict) -> dict:
    """Add derived fields to the event."""
    if not event:
        return None
    event["event_date"] = event["timestamp"][:10]
    event["event_category"] = categorize_event(event.get("event_type", ""))
    return event

def categorize_event(event_type: str) -> str:
    if event_type.startswith("page_"):
        return "navigation"
    elif event_type.startswith("purchase_"):
        return "conversion"
    return "interaction"

def run_streaming_pipeline():
    options = PipelineOptions()
    options.view_as(StandardOptions).streaming = True

    gcp_options = options.view_as(GoogleCloudOptions)
    gcp_options.project = "my-project"
    gcp_options.region = "europe-west4"
    gcp_options.job_name = "user-events-pipeline"
    gcp_options.staging_location = "gs://my-project-dataflow/staging"
    gcp_options.temp_location = "gs://my-project-dataflow/temp"

    with beam.Pipeline(options=options) as p:
        events = (
            p
            | "ReadPubSub" >> beam.io.ReadFromPubSub(
                topic="projects/my-project/topics/user-events",
                with_attributes=False,
            )
            | "DecodeBytes" >> beam.Map(lambda x: x.decode("utf-8"))
            | "ParseEvents" >> beam.Map(parse_event)
            | "FilterNone" >> beam.Filter(lambda x: x is not None)
            | "EnrichEvents" >> beam.Map(enrich_event)
        )

        # Write to BigQuery
        events | "WriteToBigQuery" >> beam.io.WriteToBigQuery(
            table="my-project:analytics.user_events",
            schema="SCHEMA_AUTODETECT",
            write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
            create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
        )

if __name__ == "__main__":
    run_streaming_pipeline()

Windowing: Aggregating Streaming Data

Windowing lets you aggregate events over time periods. Without windowing, GroupByKey in streaming mode would wait forever for all events to arrive.

import apache_beam as beam
from apache_beam.transforms.window import FixedWindows, SlidingWindows, Sessions
from datetime import timedelta

def run_windowed_aggregation():
    options = PipelineOptions()
    options.view_as(StandardOptions).streaming = True

    with beam.Pipeline(options=options) as p:
        events = (
            p
            | "ReadPubSub" >> beam.io.ReadFromPubSub(
                topic="projects/my-project/topics/user-events",
                timestamp_attribute="timestamp",  # Use event timestamp, not processing time
            )
            | "ParseEvents" >> beam.Map(lambda x: json.loads(x.decode("utf-8")))
        )

        # Count events per event_type in 5-minute fixed windows
        event_counts = (
            events
            | "ExtractType" >> beam.Map(lambda e: (e["event_type"], 1))
            | "FixedWindow5min" >> beam.WindowInto(FixedWindows(5 * 60))  # 5 minutes
            | "CountPerType" >> beam.CombinePerKey(sum)
            | "AddWindowInfo" >> beam.ParDo(AddWindowTimestampFn())
        )

        # Sliding window: events per minute, updated every 30 seconds
        sliding_counts = (
            events
            | "SlidingWindow" >> beam.WindowInto(
                SlidingWindows(60, 30)  # 60-second window, slide every 30s
            )
            | "CountAll" >> beam.combiners.Count.Globally()
        )

        # Session windows: group events by user within activity sessions
        user_sessions = (
            events
            | "KeyByUser" >> beam.Map(lambda e: (e["user_id"], e))
            | "SessionWindows" >> beam.WindowInto(
                Sessions(gap_size=30 * 60)  # 30-minute inactivity gap
            )
            | "GroupByUser" >> beam.GroupByKey()
            | "ComputeSessionMetrics" >> beam.Map(compute_session_metrics)
        )

class AddWindowTimestampFn(beam.DoFn):
    def process(self, element, window=beam.DoFn.WindowParam):
        event_type, count = element
        yield {
            "event_type": event_type,
            "count": count,
            "window_start": window.start.to_utc_datetime().isoformat(),
            "window_end": window.end.to_utc_datetime().isoformat(),
        }

Handling Late Data with Triggers and Allowed Lateness

Real-world streaming data arrives late — mobile events buffered during offline periods, retry logic that delivers events out of order, etc.

from apache_beam.transforms.trigger import AfterWatermark, AfterProcessingTime, AccumulationMode, AfterCount

events_with_late_handling = (
    events
    | "Window" >> beam.WindowInto(
        FixedWindows(5 * 60),
        trigger=AfterWatermark(
            early=AfterProcessingTime(60),    # Emit early results after 60s of processing time
            late=AfterCount(10),               # Emit late results after every 10 late arrivals
        ),
        accumulation_mode=AccumulationMode.ACCUMULATING,  # Accumulate all results
        allowed_lateness=timedelta(hours=2),   # Accept events up to 2 hours late
    )
    | "GroupByKey" >> beam.GroupByKey()
)

Stateful Processing for Advanced Patterns

Stateful processing lets you maintain state across elements — useful for session tracking, rate limiting within the pipeline, and deduplication.

from apache_beam.transforms.userstate import BagStateSpec, ReadModifyWriteStateSpec, CombiningValueStateSpec
import apache_beam as beam

class DeduplicateEventsFn(beam.DoFn):
    """Emit each event_id only once within a 1-hour window."""

    SEEN_IDS = BagStateSpec("seen_ids", beam.coders.StrUtf8Coder())

    def process(self, element, seen_ids=beam.DoFn.StateParam(SEEN_IDS)):
        event_id, event = element

        # Check if we've seen this event_id
        seen = list(seen_ids.read())
        if event_id in seen:
            return  # Skip duplicate

        seen_ids.add(event_id)
        yield event

class SessionTrackerFn(beam.DoFn):
    """Track session state per user and emit session end events."""

    SESSION_COUNT = CombiningValueStateSpec("session_count", beam.transforms.combinefn.CountCombineFn())
    LAST_SEEN = ReadModifyWriteStateSpec("last_seen", beam.coders.FloatCoder())

    SESSION_TIMEOUT = 30 * 60  # 30 minutes

    def process(self, element, session_count=beam.DoFn.StateParam(SESSION_COUNT),
                last_seen=beam.DoFn.StateParam(LAST_SEEN), timestamp=beam.DoFn.TimestampParam):
        user_id, event = element
        current_time = float(timestamp)

        last_time = last_seen.read() or 0

        # Check if this is a new session
        if current_time - last_time > self.SESSION_TIMEOUT:
            session_count.add(1)

        last_seen.write(current_time)

        yield {
            "user_id": user_id,
            "event": event,
            "session_number": session_count.read(),
        }

Dataflow Flex Templates

Flex Templates package your pipeline as a Docker container, making it easy to launch pipelines with different parameters without redeploying code:

# Build the Flex Template container
gcloud builds submit   --tag=europe-west4-docker.pkg.dev/my-project/templates/user-events:latest   --file=Dockerfile

# Create the template spec
gcloud dataflow flex-template build   gs://my-project-dataflow/templates/user-events.json   --image=europe-west4-docker.pkg.dev/my-project/templates/user-events:latest   --sdk-language=PYTHON   --metadata-file=metadata.json

# metadata.json defines pipeline parameters
cat > metadata.json << 'EOF'
{
  "name": "User Events Pipeline",
  "description": "Streams user events from Pub/Sub to BigQuery",
  "parameters": [
    {
      "name": "input_topic",
      "label": "Input Pub/Sub topic",
      "helpText": "projects/PROJECT/topics/TOPIC",
      "isOptional": false
    },
    {
      "name": "output_table",
      "label": "Output BigQuery table",
      "helpText": "PROJECT:DATASET.TABLE",
      "isOptional": false
    }
  ]
}
EOF

# Launch the template with parameters
gcloud dataflow flex-template run user-events-job   --template-file-gcs-location=gs://my-project-dataflow/templates/user-events.json   --region=europe-west4   --parameters=input_topic=projects/my-project/topics/user-events,output_table=my-project:analytics.user_events   --max-workers=20   --worker-machine-type=n1-standard-4

Pipeline Monitoring and Operations

# Monitor system lag (time from oldest unprocessed message to now)
gcloud monitoring metrics list   --filter="metric.type=dataflow.googleapis.com/job/system_lag"

# Get pipeline execution graph and metrics
gcloud dataflow jobs describe JOB_ID   --region=europe-west4   --format=json | jq '.currentState, .currentStateTime'

# Drain a streaming job gracefully (processes buffered data before stopping)
gcloud dataflow jobs drain JOB_ID   --region=europe-west4

# Cancel a job immediately (for non-production or emergencies)
gcloud dataflow jobs cancel JOB_ID   --region=europe-west4

Create a monitoring dashboard for your Dataflow pipeline:

# Alert when system lag exceeds 5 minutes
gcloud alpha monitoring policies create   --display-name="Dataflow Pipeline Lag"   --condition-filter='resource.type="dataflow_job" AND metric.type="dataflow.googleapis.com/job/system_lag"'   --condition-threshold-value=300   --condition-threshold-comparison=COMPARISON_GT   --condition-duration=300s   --notification-channels=$CHANNEL_ID

For connecting your Dataflow pipeline to BigQuery for analytics, see our BigQuery streaming guide. For the Pub/Sub foundation, see our Cloud Run event-driven microservices guide.