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.