Streaming: Spark Structured Streaming

The nightly revenue report from the first post in this track now finishes in twenty minutes — a win. But product's "dashboard that updates every minute" request is still open, and yesterday's incident made it urgent: a pricing bug overcharged customers for three hours before anyone noticed, because the batch job only runs at midnight. "We need to see what's happening now," the on-call engineer says. "Not what happened yesterday."

That's the streaming question — "what's happening?" — and this post answers it with the engine you already know. Spark Structured Streaming takes the DataFrame API from posts 3–4 and points it at data that never ends: instead of learning a separate streaming API, you write what looks like a batch query, and the engine executes it incrementally as new rows arrive. By the end you'll understand the mental model, the one interview concept that matters (event time vs processing time), watermarks, windows, triggers, output modes, and — honestly — what "exactly-once" actually requires.

Lab honesty, up front. This is the one post in the track whose lab cannot run in this sandbox: Structured Streaming is built on Spark's SQL/DataFrame engine, and that engine cannot start here — its internal Netty transport fails with a "too large frame" error (verified by direct investigation in post 1; plain RDD operations work, anything touching Spark SQL does not). Every code listing below therefore carries an explicit [NOT RUN] label: each listing is complete and runnable on a normal machine and describes standard, documented Spark behavior, but nothing was executed here and no output below was produced here. Honest unrunnable code beats an invented streaming run every time.

The mental model: the unbounded table

Structured Streaming's central idea fits in one sentence: a stream is a table that never stops growing, and your query is the same query you'd write against a batch table. The API makes this literal — reads and writes come in pairs:

  • spark.read reads a bounded table; spark.readStream reads an unbounded one. Both return a Dataset<Row>.
  • df.write writes a result once; df.writeStream keeps writing results as they update. The transformations in between are identical.

Those transformations — filter, select, groupBy, agg — don't know or care which kind of table they're running against. What changes is the execution: instead of scanning a finished dataset once, the engine re-runs your query over each new batch of arriving rows (micro-batches) and incrementally updates a result table. Your windowed count of clicks per URL isn't a different kind of query from the batch version — it's the same query, executed repeatedly as the input table grows.

One query, two executions Unbounded input Kafka topic: clicks new rows arrive forever 12:03:11 /pricing 12:03:14 /docs 12:03:19 /pricing ⋮ Streaming query same operators as batch: filter · groupBy · agg the query doesn't know it's streaming — the engine runs it incrementally Result table updated as data arrives 12:00–12:10 /pricing → 148 /docs → 96 updated every trigger Write the query once, as batch. Each micro-batch re-executes it over the new rows only.

Two consequences fall out of this model, and they explain most of the API. First, because the input never ends, some queries are simply illegal on streams: sorting an unbounded table has no final answer, so the engine rejects it instead of hanging. Second, because the result table keeps changing, you need a policy for what gets written and when — that's what triggers and output modes are (below). The model stays honest throughout: anything that would require seeing "all" of an infinite table is either restricted or needs a watermark to bound it.

Event time vs processing time: the interview concept

Every event carries two clocks, and confusing them is the most common streaming bug:

  • Event time is when the thing happened — the timestamp inside the payload, set by the producer. A click made at 12:03:11 has event time 12:03:11 even if it reaches your engine at 12:09.
  • Processing time is when your engine noticed — the machine's clock at the moment the record is processed. That same click, processed at 12:09:44, has processing time 12:09:44.

Principle: event time is when it happened; processing time is when you noticed. Design windows on the first; debug with the second.

In a perfect network the two agree and the distinction feels academic. Networks aren't perfect: mobile clients go offline and sync an hour later, retries delay and duplicate deliveries, upstream jobs backfill yesterday's data today. When that happens, processing-time windows silently give the wrong answer — a click from 12:03 lands in the 12:09 window because that's when it arrived, and reprocessing the same log tomorrow produces different results. Event-time windows are deterministic: the same input data always yields the same windows, regardless of when the engine happened to see each row. That's the interview answer in one line: event time makes results reproducible; processing time makes them an accident of arrival order.

Structured Streaming's event-time support is built around a timestamp column in your data plus two declarations you'll meet next: a watermark (how long to wait for late data) and a window (how to group event times). The engine tracks event time per row; processing time is simply what you get when you never declare an event-time column.

Watermarks: how long to wait for late data

Late data is inevitable, but you can't wait forever — a window that never closes holds its state in memory indefinitely. A watermark is the engine's answer: a threshold on event time that says "we now believe no events older than this will arrive." You declare it as a delay, not an absolute time:

[NOT RUN] Spark's SQL engine can't run in this lab (Netty transport failure); this code is complete and runnable on a normal machine — it describes standard documented behavior.

Dataset<Row> perWindow = events
        .withWatermark("event_time", "10 minutes")   // wait up to 10 min for late data
        .groupBy(window(col("event_time"), "10 minutes"), col("url"))
        .count();

Behind that one line, the engine maintains watermark = (maximum event time seen so far) − 10 minutes. The watermark only moves forward, and only when rows carrying newer event times arrive. It does two jobs, and both matter:

  1. It decides when "late" is too late. A row whose event time is older than the watermark is dropped — the engine has moved on. With a 10-minute watermark, an event timestamped 12:05 that arrives after the watermark has passed 12:05 never affects any result.
  2. It bounds the engine's state. Once the watermark passes a window's end, that window's intermediate aggregation state is discarded. Without a watermark, the engine would have to keep every window it ever saw in memory — unbounded input would mean unbounded state, and the query would eventually fall over.

Two details worth knowing precisely. First, a watermark only does anything on queries that aggregate over event time — declaring one on a plain filter query changes nothing. Second, the delay is a business decision, not a tuning knob: "10 minutes" means "we accept losing data more than 10 minutes late in exchange for closing windows promptly." Shorter watermarks mean fresher results and less state; longer ones mean fewer dropped events. No universally right value exists — it comes from how late your producers can realistically be.

Windows: tumbling, sliding, session

A window groups event times into buckets. Spark offers three shapes; the diagram shows the same events grouped all three ways:

Windows over event time on-time late → dropped watermark 12:21 Tumbling 10-min fixed 12:00–12:10 12:10–12:20 12:20–12:30 12:30–12:40 late Sliding 10-min win 5-min slide Session 5-min gap 12:03–12:14 session 8-min gap > 5-min inactivity → new session 12:00 12:10 12:20 12:30 12:40 Same events, three groupings. The watermark decides when "late" is too late — and lets the engine drop old state.
  • Tumbling windows are fixed-size and non-overlapping: window(col("event_time"), "10 minutes") puts each event in exactly one 10-minute bucket (12:00–12:10, 12:10–12:20, …). The workhorse for "per-10-minute counts" dashboards.
  • Sliding windows are fixed-size and overlapping: window(col("event_time"), "10 minutes", "5 minutes") starts a new 10-minute bucket every 5 minutes. Each event belongs to several windows — the standard shape for moving averages ("average latency over the last 10 minutes, recomputed every 5").
  • Session windows are activity-based: session_window(col("event_time"), "5 minutes") groups events separated by less than 5 minutes of inactivity into one session and starts a new session after a longer gap. Sessions have dynamic lengths — the engine merges them as events arrive — which makes them the right shape for "user session" analytics and the wrong shape for fixed reporting periods.

Notice the late event in the top row: timestamped 12:05 but arriving after the watermark has already passed 12:05. The engine drops it — the 12:00–12:10 window already closed. That stings the first time you see it, but the alternative is worse: without the watermark's cutoff, the engine could never finalize any window, and "per-10-minute counts" would be a query that never produces an answer.

Lab: the full pipeline

[NOT RUN] Spark's SQL engine can't run in this lab (Netty transport failure); this code is complete and runnable on a normal machine — it describes standard documented behavior.

Here is the whole thing: Kafka in, watermarked event-time aggregation, console out. It assumes a Kafka broker on localhost:9092 with a clicks topic carrying JSON payloads like {"user_id":"u42","event_time":"2026-10-06T12:03:11","url":"/pricing"} (you built exactly this setup in the Data & Messaging track). The Spark Kafka connector ships separately — pass --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3 to spark-submit — and run it on JDK 17, Spark 3.5's supported runtime (post 1). Compile the class against your Spark installation's jars ($SPARK_HOME/jars/* on the classpath).

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.streaming.StreamingQuery;
import org.apache.spark.sql.streaming.Trigger;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.StructType;

import static org.apache.spark.sql.functions.*;

public class ClicksPerWindow {
    public static void main(String[] args) throws Exception {
        SparkSession spark = SparkSession.builder()
                .appName("ClicksPerWindow")
                .master("local[*]")
                .getOrCreate();
        spark.sparkContext().setLogLevel("WARN");

        // 1. SOURCE: an unbounded table over the Kafka topic "clicks".
        //    readStream returns immediately; nothing is read until the query starts.
        Dataset<Row> kafka = spark.readStream()
                .format("kafka")
                .option("kafka.bootstrap.servers", "localhost:9092")
                .option("subscribe", "clicks")
                .option("startingOffsets", "latest")
                .load();

        // 2. PARSE: Kafka delivers key/value as binary. The value is JSON;
        //    from_json promotes it to typed columns, including a real timestamp.
        StructType schema = new StructType()
                .add("user_id", DataTypes.StringType)
                .add("event_time", DataTypes.TimestampType)
                .add("url", DataTypes.StringType);

        Dataset<Row> events = kafka
                .selectExpr("CAST(value AS STRING) AS json")
                .select(from_json(col("json"), schema).as("e"))
                .select("e.*");

        // 3. EVENT-TIME AGGREGATION: 10-minute tumbling windows on event_time.
        //    The watermark waits up to 10 minutes for late data, then closes
        //    each window — which is also what bounds the engine's state.
        Dataset<Row> perWindow = events
                .withWatermark("event_time", "10 minutes")
                .groupBy(window(col("event_time"), "10 minutes"), col("url"))
                .count();

        // 4. SINK: append mode is legal here because the watermark marks each
        //    window's result final. The checkpoint makes the query resumable:
        //    Kafka offsets and aggregation state survive a restart.
        StreamingQuery query = perWindow.writeStream()
                .outputMode("append")
                .format("console")
                .option("truncate", "false")
                .option("checkpointLocation", "/tmp/ss-checkpoint/clicks-per-window")
                .trigger(Trigger.ProcessingTime("1 minute"))
                .start();

        query.awaitTermination();
    }
}

Run it with:

$ spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3 \
    --class ClicksPerWindow --master "local[*]" clicks-per-window.jar

[NOT RUN] Spark's SQL engine can't run in this lab (Netty transport failure); the shell session below is the standard documented way to feed the pipeline above on a normal machine — it was not executed here.

To feed the pipeline test data, use the Kafka console producer you already know — type JSON lines, then Ctrl+C when done:

$ kafka-console-producer.sh --topic clicks --bootstrap-server localhost:9092
{"user_id":"u42","event_time":"2026-10-06T12:03:11","url":"/pricing"}
{"user_id":"u17","event_time":"2026-10-06T12:07:52","url":"/docs"}
{"user_id":"u42","event_time":"2026-10-06T12:11:03","url":"/pricing"}

What this pipeline does, beat by beat: readStream returns immediately with a DataFrame representing the unbounded clicks table — nothing is read yet. The from_json parse turns Kafka's binary value into typed columns, promoting the payload's event_time to a real timestamp. The watermark plus tumbling window groups rows into 10-minute buckets per URL, and because the watermark marks each window final, append mode is legal: each window's count is emitted once, when it closes. The trigger fires every minute; the checkpoint directory records Kafka offsets and aggregation state, so killing and restarting the query resumes where it left off instead of reprocessing from scratch.

Triggers and output modes: what gets written, and when

A streaming query produces a result table that keeps changing. Two settings control what the sink sees:

Triggers control when the engine runs:

  • Default (no trigger set): start the next micro-batch as soon as the previous one finishes and new data is available — run as fast as the data arrives.
  • Fixed interval: Trigger.ProcessingTime("1 minute") — process whatever arrived during each interval. Longer intervals mean larger, more efficient batches at higher latency; shorter intervals mean the reverse. The micro-batch is a latency floor by design — this is the honest architectural difference from Flink, covered below.
  • Available-now: Trigger.AvailableNow() — process everything currently available, in one or more micro-batches, then stop. This is the backfill switch: it runs the streaming query as a batch job over what's already there.
  • Continuous exists as an experimental option aimed at very low latency, with restricted operators and weaker guarantees — not where you start.

Output modes control what gets written each time:

ModeWhat the sink receivesWhen it's legal
Append (default)Only new result rows — never revisions of rows already writtenQueries without aggregations, or event-time aggregations with a watermark (the watermark is what makes a window's result final)
UpdateNew rows plus changed rows (never deletions)Aggregations — the sink must handle a row's value changing
CompleteThe entire result table, every triggerAggregations — and only when the result table stays small enough to rewrite in full

The guarantee each mode gives the sink is the point of the table: append promises a row is written once and never revised (which is why it needs the watermark's proof of finality); update may revise rows and therefore needs an upsert-capable sink; complete promises simplicity at the cost of rewriting everything. Choosing append for a sink that can't absorb revisions, or complete for a result table with millions of groups, are the two classic mistakes — both come from picking the mode before thinking about what the sink can handle.

Exactly-once, honestly

"Exactly-once" is the most oversold phrase in streaming. Here is the honest boundary: Structured Streaming gives you exactly-once within its engine; the sink decides the end-to-end story.

Within the engine, the mechanism is checkpointing. The checkpoint directory (that checkpointLocation option) durably records two things: how far the source was read (Kafka offsets) and the versioned aggregation state. If the query dies mid-micro-batch, the restart replays from the last checkpoint — and because state updates are versioned against the batch, the replay doesn't double-count. That half the engine genuinely handles.

But the rows the engine emitted before dying may already be sitting in your sink. Whether they appear once or twice downstream depends on three things, all outside the engine's control:

  1. A replayable source. Kafka works because offsets let the engine re-read exactly the failed batch. A source that can't replay (a raw socket, say) can't offer better than at-least-once no matter what the engine does.
  2. A checkpoint location that survives the failure. No checkpoint, no resume point — a restarted query reprocesses from wherever the source starts it. On a real cluster this means durable storage (HDFS, S3), not /tmp on one machine.
  3. An idempotent or transactional sink. This is where end-to-end exactly-once is won or lost.

The sink patterns, from strongest to weakest:

PatternHow it stays exactly-onceExample
Transactional sinkResults and offsets commit atomically — a failed batch leaves no traceDelta Lake's streaming sink (tracks written batches in its transaction log); Kafka-to-Kafka with a transactional producer
Idempotent upsert by keyRe-running a batch overwrites identical rows instead of duplicating them — the natural key makes retries harmlessforeachBatch writing INSERT … ON CONFLICT DO UPDATE keyed by (window, url) — listing below
Plain appendNothing prevents duplicates on replayFiles, console, any sink without keys or transactions — honestly at-least-once

The Delta Lake case is a one-line change to the lab pipeline — same query, exactly-once sink:

[NOT RUN] Spark's SQL engine can't run in this lab (Netty transport failure); this code is complete and runnable on a normal machine — it describes standard documented behavior.

StreamingQuery query = perWindow.writeStream()
        .outputMode("append")
        .format("delta")
        .option("checkpointLocation", "/tmp/ss-checkpoint/clicks-delta")
        .start("/data/delta/clicks_per_window");

And the general-purpose escape hatch is foreachBatch: per micro-batch, the engine hands your function a plain batch DataFrame and a batch id, and you write it with arbitrary code. The contract is explicit and worth memorizing: after a failure, the engine may call your function more than once for the same batch id — so exactly-once becomes your code's responsibility. The standard answer is an upsert keyed on the data's natural key, so a re-run is an overwrite, not a duplicate. First the table (Postgres syntax — other databases spell the upsert as MERGE):

[NOT RUN] Spark's SQL engine can't run in this lab (Netty transport failure); this code is complete and runnable on a normal machine — it describes standard documented behavior.

CREATE TABLE clicks_per_window (
    window_start TIMESTAMP,
    window_end   TIMESTAMP,
    url          TEXT,
    cnt          BIGINT,
    PRIMARY KEY (window_start, url)
);

[NOT RUN] Spark's SQL engine can't run in this lab (Netty transport failure); this code is complete and runnable on a normal machine — it describes standard documented behavior.

// Same 'perWindow' query as the lab listing; only the sink changes.
// Imports needed beyond the lab listing: java.sql.Connection,
// java.sql.DriverManager, java.sql.PreparedStatement.
private static final String JDBC_URL = "jdbc:postgresql://localhost:5432/analytics";
private static final String DB_USER = "analytics";
private static final String DB_PASS = "change-me";   // fill in your own

StreamingQuery query = perWindow.writeStream()
        .outputMode("append")
        .foreachBatch((batchDF, batchId) -> {
            // HONEST CONTRACT: this function may run MORE THAN ONCE for the
            // same batchId after a failure. Exactly-once is therefore this
            // code's job: the write below is an UPSERT keyed by
            // (window_start, url), so a re-run overwrites identical rows
            // instead of duplicating them.
            Dataset<Row> flat = batchDF.select(
                    col("window.start").as("window_start"),
                    col("window.end").as("window_end"),
                    col("url"),
                    col("count").as("cnt"));

            flat.foreachPartition(rows -> {
                Connection conn = DriverManager.getConnection(JDBC_URL, DB_USER, DB_PASS);
                conn.setAutoCommit(false);
                PreparedStatement ps = conn.prepareStatement(
                        "INSERT INTO clicks_per_window(window_start, window_end, url, cnt) " +
                        "VALUES (?,?,?,?) " +
                        "ON CONFLICT (window_start, url) DO UPDATE SET " +
                        "cnt = EXCLUDED.cnt, window_end = EXCLUDED.window_end");
                while (rows.hasNext()) {
                    Row r = rows.next();
                    ps.setTimestamp(1, r.getAs("window_start"));
                    ps.setTimestamp(2, r.getAs("window_end"));
                    ps.setString(3, r.getAs("url"));
                    ps.setLong(4, r.getAs("cnt"));
                    ps.addBatch();
                }
                ps.executeBatch();
                conn.commit();
                conn.close();
            });
        })
        .option("checkpointLocation", "/tmp/ss-checkpoint/clicks-jdbc")
        .start();

query.awaitTermination();

Read the contract in that listing once more: the batchId is deterministic across re-runs, the SQL is Postgres's upsert (ON CONFLICT), and re-running any batch sets the same values on the same keys. The engine's replay guarantee plus the sink's idempotence is what "exactly-once" actually means in practice — two halves, and you own the second one.

Structured Streaming vs Flink: the rule of thumb, applied

Post 1 gave the rule of thumb; this post is where you cash it in: if seconds-to-minutes latency is acceptable and Spark is already your platform, Structured Streaming is usually enough; when consistently low latency, sophisticated event-time processing, and deeply stateful streaming are central requirements, Flink deserves serious consideration. Now you can see why it's shaped that way:

  • Latency model. Structured Streaming is a micro-batch engine: the trigger interval is a latency floor by design (with an experimental continuous mode carrying restricted operators). Flink processes event-by-event, which is why it can sit at consistently lower latencies — but you only feel that difference when your use case actually needs it.
  • Model simplicity. The unbounded-table model means your Spark SQL knowledge transfers directly — the same groupBy from post 4 works on a stream. Flink's DataStream API is a separate model to learn, with more knobs for state, timers, and process functions.
  • Concepts transfer. Event time, watermarks, windows, exactly-once sinks — everything in this post is Flink vocabulary too. Learning streaming on Structured Streaming is not a dead end if you later need Flink; it's the prerequisite.

And the second half of the rule: don't introduce a second streaming platform because a benchmark said it's faster. Two streaming systems means two operational skill sets, two failure modes, and two exactly-once stories to get right. Most teams never outgrow the first half of the sentence — and the teams that do will know it from their latency graphs, not from a blog post.

Cheat sheet

  • A stream is an unbounded table; the query is the same as batch — readStream/writeStream mirror read/write, and the engine executes incrementally.
  • Event time is when it happened; processing time is when you noticed. Event time makes results reproducible; processing time makes them an accident of arrival order.
  • A watermark is "how long to wait for late data," declared as a delay: watermark = max event time seen − delay. It drops too-late rows and bounds the engine's state — both jobs matter.
  • Tumbling windows: fixed, non-overlapping (window(col, "10 minutes")). Sliding: fixed, overlapping (window(col, "10 minutes", "5 minutes")). Session: gap-based and dynamic (session_window(col, "5 minutes")).
  • Triggers set when (default = as fast as data arrives; fixed interval; AvailableNow for backfill; continuous is experimental). Output modes set what: append (rows written once — needs watermark finality), update (changed rows — needs an upsert-capable sink), complete (whole table — only for small results).
  • Exactly-once recipe: replayable source (Kafka offsets) + durable checkpoint location + idempotent/transactional sink. The engine handles its half via checkpointing; foreachBatch hands the sink half to you — it may re-run, so upsert by key.
  • Rule of thumb, not a law: seconds-to-minutes latency on Spark already → Structured Streaming is usually enough; consistently low latency with deep stateful streaming → evaluate Flink. The concepts transfer either way.
  • Kafka connector ships separately: --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3; Spark 3.5 wants JDK 8/11/17.

What's next

Your streaming query now emits per-window counts — but where do those results live? The console sink was a teaching aid; production results land in a store built for the access pattern: high-write, key-scoped reads, replicated across data centers. That's the wide-column world, and it starts with the system that defined it.

The next post, Cassandra: Wide-Column at Scale, goes hands-on with query-first data modeling, tunable consistency, and CQL from Java — the serving layer your streaming results will eventually call home.

Field check before you move on: (1) On your own machine, run the pipeline with a 5-minute window and a 2-minute watermark; produce one event 3 minutes late and confirm whether it's counted or dropped — then explain why, using the watermark formula. (2) Switch the sink to update mode with an AvailableNow trigger and describe what changed about when results appear. (3) Kill the query mid-run, restart it, and inspect the checkpoint directory — what did the engine replay, and what did it not need to recompute? (4) Take one streaming workload from your own work and write a single sentence: event time or processing time — and what breaks if you choose wrong.

Continue: Java Learning Roadmap 2026

Comments

Popular posts from this blog

JSP Servlet Interview Questions For Freshers Series 1

Java Banking Finance Services and Insurance (BFSI) domain interview questions

Java program to check even or odd number