TL;DR / Key Takeaways
Structured Streaming checkpoints traditionally identify sources by position, so changing the source list can require a fresh checkpoint and discarded state.
Source evolution tracks offsets by stable
.name()values instead.On-demand state repartitioning (Public Preview, RocksDB) can resize state partitions on restart.
This Ethereum archive-to-live test preserved one checkpoint and query id through source addition, source removal, and state repartitioning without rereading archive rows or losing keys.
Name every source when creating the checkpoint, and plan for tombstoned names and retained state.
Problem: We hope for permanent sources, and life give us transient ones
Think about what a streaming job is actually for. It maintains a metric: orders per region per minute, active sessions per tenant, fraud score per card, blocks per proposer per window. That metric has consumers who expect it to be continuous. It has a state store carrying every open window and every in-flight key. It is the thing the business depends on.
Now think about what feeds it. A Kafka cluster that will be migrated next quarter. A vendor feed on a contract that ends in March. A region launching in six months. An S3 archive that exists only until the live feed catches up. A tenant who churns. None of these are permanent, and everyone knows it on the day the job is written.
Yet the checkpoint treats the source list as permanent. Structured Streaming identifies sources by position—0, 1, 2—and the offset log is positional too. Add, remove, or reorder a source and the query refuses to start; for a stateful job, the only fix was a fresh checkpoint—discarding state and replaying history while consumers see a gap or a double count.
Teams worked around that in three ways: replay inside a maintenance window (duplicates and downtime), one job per source (the metric assembled downstream from N partial metrics), or a bus in the middle (an extra system /table whose only job is making the source list look permanent).
Each workaround compensates for the engine; source evolution removes the need—the job and its state stay put, and sources become what they always were: replaceable.
Each of those is an architecture decision made to compensate for the engine. Source evolution removes the need for it: the job and its state stay put, and the sources become what they always were, replaceable.
What changed: named sources and resizable state
With source evolution enabled, every streaming source carries a stable, user-defined name via .name(), and the checkpoint tracks offsets by that name instead of by position. Three operations become safe across restarts on the same checkpoint:
Reorder sources. Each one resumes from its own last committed offset.
Add a source. The new source starts from its beginning; existing sources continue where they stopped.
Remove a source. Its name is tombstoned in the checkpoint and cannot be reused, which protects you from silently rewinding a feed.
There is a second thing a checkpoint used to freeze: the number of state partitions, set by spark.sql.shuffle.partitions on the day the checkpoint was born. A backfill wants a wide cluster and many partitions; a live feed of a few records per second wants a small node and a handful. On-demand state repartitioning (Public Preview, RocksDB) unfreezes that too. Together, the two features cover the whole lifecycle: change what feeds the query, then right-size what runs it, without ever starting over.
The proof: one checkpoint through a full ingestion lifecycle
We wanted a workload where losing state would hurt, on data anyone can reproduce: AWS publishes the entire Ethereum chain as date-partitioned Parquet in a public S3 bucket, and public JSON-RPC endpoints serve the live tip with no signup. The query counts blocks per 5-minute window, per source and proposer (miner), with a 2-day watermark, update mode, a foreachBatch MERGE into Delta, and RocksDB state—run through five steps on one checkpoint, as tasks of one Databricks job, scored from the run's own StreamingQueryProgress

StreamingQueryProgress records, unedited. Red: keys in RocksDB after each batch. Blue: state store partitions. Four phases, one query id.Step 1: name the sources
For a newly created compatible query, the following two configs plus a stable, unique .name() on every source enable supported source-list evolution. This cannot be enabled retroactively on an existing checkpoint; removed names are permanently tombstoned; and support requires Databricks Runtime 18.2 or later.
spark.conf.set("spark.sql.streaming.queryEvolution.enableSourceEvolution", "true")
spark.conf.set("spark.sql.streaming.offsetLog.formatVersion", "2")
archive = (spark.readStream.format("parquet").schema(BLOCK_SCHEMA)
.option("maxFilesPerTrigger", 365)
.name("s3_backfill")
.load(landing))
Naming requires a fresh checkpoint, and it is one-way. You cannot enable source evolution on a checkpoint created without it, and you cannot turn it off later. The practical consequence: name every source on day one, including the ones you are certain will never change.
If a fresh query fails with NAMED_SOURCES_REQUIRE_OFFSET_LOG_V2, set spark.sql.streaming.offsetLog.formatVersion to 2. Note that the message text refers to the key as offsetLog.version; the key that takes effect on Databricks Runtime 18.3 is formatVersion.
Step 2: backfill the archive on a cluster that turns itself off
Backfill is the expensive part: We used 8 workers, Spark’s default 200 shuffle partitions, and trigger(availableNow=True) to read everything and stop.
(events.writeStream
.option("checkpointLocation", checkpoint)
.trigger(availableNow=True) # drain, then stop paying
.foreachBatch(upsert_into_delta)
.start())
298394343088892: the backfill receipt for the checkpoint the rest of this post continues from.The query processed 26,086,569 Ethereum blocks from 4,080 files in 13 micro-batches, leaving 3,154 keys in state. The checkpoint now retained 200 state partitions.. That number is now baked into the checkpoint. We come back to it in Step 5.
Step 3: add the live feed, and watch the archive replay nothing
This is the cutover that used to cost a fresh checkpoint. The live feed is a Python data source over a public Ethereum RPC endpoint, configured to start one block after the archive’s last. It runs on a single small node, on the same checkpoint, with the archive source still present.
The important thing is what does not change. Both sources project to the same seven columns, so the union, the watermark, the aggregation and the sink are exactly the code that ran in Step 2. The only edit is one more entry in the list of sources:
archive = (spark.readStream.format("parquet").schema(BLOCK_SCHEMA)
.option("maxFilesPerTrigger", 365)
.name("s3_backfill") # unchanged
.load(landing))
live = (spark.readStream.format("eth_rpc") # Python DataSource
.option("start_block", archive_max_block + 1) # one past the archive
.name("eth_live") # NEW named source
.load())
# Same projection for every source, so the union is by name.
events = (archive.select(*BLOCK_COLUMNS).withColumn("source", lit("s3_backfill"))
.unionByName(
live.select(*BLOCK_COLUMNS).withColumn("source", lit("eth_live"))))
# The metric. Unchanged from Step 2.
counts = (events
.withWatermark("block_ts", "2 days")
.groupBy(window("block_ts", "5 minutes"), "source", "miner")
.count())
(counts.writeStream
.outputMode("update")
.option("checkpointLocation", checkpoint) # SAME checkpoint as Step 2
.foreachBatch(upsert_into_delta) # MERGE on (window, source, miner)
.start())Spark reads the offset log by name, finds s3_backfill at its last committed offset, finds no entry for eth_live, and starts the new source from its beginning. The RocksDB state for every open window carries straight over, because the stateful operator did not change.
To see what happened on restart, we flattened the query’s own StreamingQueryProgress events (query.recentProgress, or a StreamingQueryListener) into one row per source per microbatch. Each event carries a sources[] array with a SourceProgress per input: its description, numInputRows, startOffset and endOffset.

699608591493844: one row per source per batch, flattened from StreamingQueryProgress.sources[]. Batch ids continue from 13, where the backfill stopped.Three things to read off this table:
The batch ids start at 13, not 0. Step 2 committed batches 0 through 12. A fresh checkpoint would have restarted the count; this one continued it.
The
FileStreamSourcerows (srcIdx 0) read 0 input rows in every batch, withstartOffsetandendOffsetboth held at{"logOffset": 11}. The archive resumed exactly where it stopped and re-read nothing.The new source (
srcIdx 1) starts fromnullin batch 13, then walks forward 20 blocks per batch from 26,086,590. It had no committed offset, so it began at its configured start block.
The new source started from its beginning; the old one resumed exactly where it stopped, and no history was re-read. Nothing about the aggregation state was touched.
Why is the
namecolumn empty when every source has a.name()?
.name()is a checkpoint identity, not a runtime label: on Databricks Runtime 18.3,SourceProgressdescribes sources by class and description. Names live in the checkpoint’soffsets/directory, in V2 entries keyed by source name (Step 4 reads one); verify bysrcIdx, offsets, and row counts.
Step 4: Retire the source or keep it running, depends on your use-case
Once the live feed has been running for a while, the archive has nothing left to contribute. Retiring it is not a migration; it is a restart with a shorter list. Stop the query, edit the same pipeline code and delete s3_backfill code, and start it against the same checkpoint. Nothing else is edited:
query.stop()
live = (spark.readStream.format("eth_rpc")
.option("start_block", archive_max_block + 1) # ignored: offset comes from the checkpoint
.name("eth_live") # unchanged
.load())
# No union any more. Same projection, same metric, same sink.
events = live.select(*BLOCK_COLUMNS).withColumn("source", lit("eth_live"))
counts = (events
.withWatermark("block_ts", "2 days")
.groupBy(window("block_ts", "5 minutes"), "source", "miner")
.count())
(counts.writeStream
.outputMode("update")
.option("checkpointLocation", checkpoint) # SAME checkpoint
.foreachBatch(upsert_into_delta)
.start())On restart Spark compares the named sources in the plan with the named sources in the newest offset-log entry. eth_live is in both, so it resumes from its committed offset (the start_block option is only honoured the first time a name is seen). s3_backfill is in the log but not in the plan, so Spark records its departure. The offset log keeps the name, marked with a tombstone, and the live source carries on:

699608591493844. Top: the newest offset-log entry after the removal. Bottom: the state store read back against Delta.Two behaviours are worth noticing. First, removing a source does not remove its state. At removal, RocksDB still held 3,060 archive keys, exactly one per archive row in Delta whose window the watermark had not yet closed. Those windows will close on their own schedule, as they should. Second, the name is gone for good. Re-attach s3_backfill and the query refuses to start:
[STREAMING_QUERY_EVOLUTION_ERROR.TOMBSTONE_SOURCE_NAME_REUSE] Cannot reuse tombstoned source names: s3_backfill. These source names were previously used and then removed from the streaming query at checkpoint location ... Reusing tombstoned source names can lead to data correctness issues. Please use different source names.You want that error. A reusable name would let a source silently rewind to an offset from two restarts ago and double-count every window it re-read.
Step 5: shrink the state to fit the live workload
With the archive gone, the query is processing 20 blocks per batch on one node, but its state is still spread across the backfill’s 200 partitions. Changing spark.sql.shuffle.partitions in the session does nothing; the checkpoint remembers 200.
# Resizes state on the next restart. shuffle.partitions cannot.
spark.conf.set("spark.sql.streaming.stateStore.partitions", "8")
StreamingQueryProgress records, unedited. Every point is a batch of exactly 20 blocks.The restart’s first progress event reports 66.8 seconds of controlBatch.REPARTITION, with no keys lost: 3,155 before, 3,157 after, 0 removed. From then on the same 20-block batches took a median 15.0 seconds of addBatch instead of 35.5, and state memory fell from 1,323 MB to 1,067 MB. Same checkpoint, same query id, a third of the partitions’ overhead gone.
Proof: two clusters, four phases, one query id
Spark derives a streaming query’s id from checkpoint metadata, so an unchanged id across restarts is direct evidence that all four phases were the same logical query resuming, not four queries that happened to write to the same table. Both clusters’ Spark UIs show it:

Nothing fell through the seams. Delta holds 26,087,209 blocks: the archive’s 26,086,569 plus 640 live blocks starting at 26,086,570, exactly one past the archive’s last.
What this means for how you design streams
The demo is a blockchain, but the principal is universal: a long-lived metric fed by sources which will come and go. Once the job no longer has to be rebuilt when its inputs change, several patterns become routine restarts instead of migration projects:
And four rules to carry into every new stream:
Assume every source is temporary, and name it on day one. Naming needs a fresh checkpoint, so the cheapest moment is before the first batch, and the sources you are most certain about are the ones you will be most surprised by.
Treat names as permanent identifiers. A removed name is tombstoned; a rename is a remove plus an add that starts from the beginning.
Removing a source keeps its state until the watermark closes it. Plan for that in memory sizing.
Remember to alter
spark.sql.streaming.stateStore.partitions: when you’re backfilling set this number high so you can backfill fast, and then remember to bring down that number when you go in live state.
Common misconceptions about streaming source lifecycle management
The lifecycle test exposes several common misconceptions; confusing identity, offsets, and retained state causes replay, double counting, or failed restarts.
Misconception: a source list is fixed for the life of its checkpoint. Correction: with source evolution from a fresh checkpoint and stable names, sources can be added, removed, or reordered across restarts without discarding state.
Misconception: adding a source causes existing sources to replay. Correction: existing named sources resume from their committed offsets; only the newly introduced name starts from its configured beginning.
Misconception: removing a source immediately deletes its state. Correction: removal tombstones the name; existing state remains until normal watermark eviction closes it.
Misconception: a removed source name can later be reused. Correction: tombstoned names cannot be reused, because doing so could rewind offsets and double-count data.
Misconception: source names appear as labels in
StreamingQueryProgress. Correction: on the tested runtime,.name()is a checkpoint identity; verify by source index, offsets, row counts, or V2 offset-log entries.Misconception: changing
spark.sql.shuffle.partitionsresizes an existing state store. Correction: the checkpoint keeps its count; usespark.sql.streaming.stateStore.partitionswith on-demand repartitioning.
Scope of this run
Runtime and compute. Databricks Runtime 18.3 on classic compute. Source evolution requires Databricks Runtime 18.2 or above; on-demand state repartitioning is in Public Preview and requires the RocksDB state store.
Structured Streaming jobs. This run used plain Structured Streaming with
foreachBatch. It did not exercise Lakeflow Spark Declarative Pipelines or real-time mode.Not a benchmark. One workload, one run. The partition-count timings are illustrative of the mechanism, not a performance claim. Things can be done in milli-seconds end to end as well.
Get started
Two configs and a .name() on each source is all it takes to decouple the job from its inputs:
spark.conf.set("spark.sql.streaming.queryEvolution.enableSourceEvolution", "true")
spark.conf.set("spark.sql.streaming.offsetLog.formatVersion", "2")
orders_us = spark.readStream.name("orders_us").table("catalog.schema.orders_us")
orders_eu = spark.readStream.name("orders_eu").table("catalog.schema.orders_eu")
all_orders = orders_us.union(orders_eu)FAQ
Q: What is source evolution in Structured Streaming?
A: Source evolution lets a new checkpoint track each source by a stable .name() rather than query-plan position. With offset-log format V2 and an otherwise compatible plan, named sources can be added, removed, or reordered across restarts.
Q: Do I need a new checkpoint to enable source evolution?
A: Yes. Source evolution cannot be enabled retroactively, and it cannot be disabled after the checkpoint is created. It requires Databricks Runtime 18.2 or later.
Q: Does adding a source replay existing sources?
A: No. Existing named sources resume from their committed offsets. Only the new source name starts from its configured beginning.
Q: What happens to state when a source is removed?
A: The source name is tombstoned, but its existing aggregation state remains until normal watermark eviction closes the associated windows.
Q: Can I change the number of state partitions without losing state?
A: With supported on-demand state repartitioning and RocksDB, set spark.sql.streaming.stateStore.partitions and restart. Changing spark.sql.shuffle.partitions does not resize state already recorded in the checkpoint.
Q: Can a tombstoned source name be reused?
A: No. Spark blocks reuse because an old committed offset could rewind the feed and double-count data. A returning feed requires a different name and therefore starts from that new name’s configured beginning.



