Spark Structured Streaming

Structured Streaming is Apache Spark's stream processing engine. Its central idea is elegant: treat a stream as an unbounded table that keeps growing, and write the same DataFrame or SQL query you'd write for batch data. Spark runs that query incrementally, processing new data as it arrives and updating results, while handling fault tolerance, state, and exactly-once output for you.

It's widely used for streaming ETL into data lakehouses, ingesting from Kafka, real-time aggregations and alerting, and change data capture pipelines. It sits alongside engines like Flink in the broader stream processing landscape. Its strengths are the unified batch and streaming API and tight integration with Spark's ecosystem, and its default micro-batch model trades some latency for throughput and simplicity.

TL;DR

Quick Example

Kafka → windowed aggregation → Delta table:

Core Concepts

The Unbounded Table Model

Each new record is a row appended to an input table. Your query defines a result table, and at each trigger Spark computes what changed and writes it to the sink. Most DataFrame operations work on streams: projections, filters, joins, aggregations, and windows. Some operations (like sorting a non-aggregated stream) aren't supported, because they'd require seeing the whole unbounded input.

Triggers and Execution Modes

availableNow is a powerful pattern: streaming semantics (checkpoints, exactly-once, incremental processing) with batch-like scheduling and cost.

Sources

Event Time and Watermarks

Events arrive out of order and late. Use the event time in the data (not arrival time) for windows. A watermark (withWatermark("event_time", "10 minutes")) tells Spark how late data can be: the watermark = max event time seen − the threshold. It lets Spark finalize windows and drop old state. Data arriving later than the watermark is discarded.

Choosing the threshold trades completeness against latency and state size.

Windows

Output Modes

State and Stateful Processing

Aggregations, deduplication, and stream-stream joins keep state in a state store (RocksDB is recommended for large state), checkpointed for recovery. Arbitrary stateful logic (sessionization, pattern detection, timers) uses applyInPandasWithState or the newer transformWithState APIs. State grows without bound unless watermarks or TTLs expire it, so always define them.

Stream-Stream Joins

Joining two streams (such as impressions with clicks) requires buffering both sides in state. Provide watermarks on both sides plus a time-range condition (clicks.ts BETWEEN impressions.ts AND impressions.ts + interval 1 hour), so Spark knows when buffered rows can be dropped. Stream-static joins (enriching with a dimension table) are simpler and common.

Checkpoints and Exactly-Once

The checkpoint directory stores source offsets, the commit log, and state snapshots. On restart, Spark resumes from the last committed batch and replays the uncommitted one. Combined with replayable sources (Kafka, files) and idempotent or transactional sinks (Delta and Iceberg commits keyed by batch ID), this gives end-to-end exactly-once results. With foreachBatch, you get the batch_id, so use it to make writes idempotent (MERGE on keys, or a record of processed batch IDs). See idempotency.

Checkpoints are tied to the query's logic: many query changes (like changing aggregation keys) require a new checkpoint.

Best Practices

Use foreachBatch for Upserts and Multiple Sinks

Monitor Streaming Health

Track input rows per second against processed rows per second, batch duration against the trigger interval, state store size, and watermark lag. query.lastProgress and streaming listeners expose these. Alert when processing falls behind.

Size State Deliberately

Use the RocksDB state store for large state, set watermarks and TTLs, and pick window sizes consciously. Unbounded state eventually causes OOMs and slow checkpoints.

Plan for Schema Changes

Upstream schema evolution can break streams. Use schema registries, tolerant parsing (a from_json with a known schema and error handling), and table formats with schema evolution. See event schema evolution.

Common Mistakes

Aggregations Without Watermarks

State for every key and window since the stream started accumulates forever. Always add a watermark for event-time aggregations and deduplication.

Deleting or Sharing Checkpoints

Deleting a checkpoint restarts from the configured starting offsets (duplicates or data loss), and two queries sharing a checkpoint corrupt each other. One checkpoint per query, stored durably.

Micro-Batches Longer Than the Trigger

If each batch takes 60 seconds with a 30-second trigger, the stream falls further behind. Scale resources, reduce per-batch work, or cap intake with maxOffsetsPerTrigger.

FAQ

Is Spark Structured Streaming real-time?

By default it uses micro-batches, with typical latencies of seconds up to about a minute, which suits most ETL and analytics. Lower-latency continuous and real-time modes exist for millisecond-scale needs. Flink is often chosen when very low latency with complex event-at-a-time processing is required.

What is a watermark in Spark streaming?

A moving threshold, based on the maximum event time seen minus an allowed lateness, that tells Spark when it can finalize event-time windows and discard old state. Events older than the watermark are dropped as too late.

How does Spark achieve exactly-once processing?

It records source offsets and state in checkpoints, processes each micro-batch deterministically, and writes to sinks that commit atomically or idempotently per batch ID. After failures it replays uncommitted batches without duplicating committed output.

Should I use Structured Streaming or Kafka Streams / Flink?

Use Structured Streaming when you're already on Spark and need lakehouse integration and a unified batch and stream codebase. Kafka Streams fits lightweight, library-embedded processing in JVM services. Flink excels at low-latency, large-state, event-at-a-time processing. See stream processing.

Related Topics

References