Forum Discussion
The Partition Count You Picked on Day 1 Doesn't Have to Follow You Forever
Every stateful Structured Streaming query carries a decision that used to feel permanent: the shuffle partition count chosen on the day its checkpoint was created. Get that number wrong, usually because the real production volume only shows up months after the default value of 200 got accepted without much thought, and the historical choice was binary: live with the wrong partition count indefinitely, or abandon the entire checkpoint (and the state accumulated in it) to recreate the query from scratch with the right value. For a job that has already processed months of user sessions or aggregation windows, "abandon the checkpoint" was never a real option.
On-demand state repartitioning, now in Public Preview on Databricks Runtime 18, exists precisely to close that gap: change the partition count of a stateful query without losing what has already been accumulated.
The mechanism: why this was never trivial
The state of a stateful query in Structured Streaming isn't a loose configuration number, it's physical data, stored in RocksDB instances inside the checkpoint, one per partition. Every key in your aggregation is distributed across partitions via hashing, so the partition count isn't just config metadata, it literally defines where each key physically lives on disk. Changing spark.sql.shuffle.partitions after the checkpoint already exists had no effect at all, because the engine kept reading the partition count recorded during the first run, and changing that number without redistributing the physical data first would break the relationship between key and partition.
On-demand state repartitioning solves this with a single, explicit operation: the query finishes the pending micro-batch, runs a redistribution that rehashes every key to the new partition count, and only then resumes normal processing. It's a controlled pause, not a background migration, and its duration is auditable in the StreamingQueryProgress event, under the durationMs.controlBatch.REPARTITION field.
My take: the detail that stands out here isn't the feature itself, it's how much it exposes a silent technical debt that probably exists in production right now. A lot of teams accept the default of 200 for spark.sql.shuffle.partitions without questioning it, because when the stateful query is first written nobody yet knows what the real volume will look like six months later. Before this feature, that day-1 decision effectively became permanent. It's worth mapping which of your oldest stateful queries have never had their partition count revisited since creation, that's your list of candidates for a free win.
Hands-on: scaling a query already in production
The requirement is Databricks Runtime 18 LTS or higher, with the RocksDB state store provider (already the default since DBR 17.3). The change itself is done by stopping the query, adjusting the configuration, and restarting with the same checkpoint:
# Original query, created with the default of 200 partitions
query = (df
.withWatermark("event_time", "10 minutes")
.groupBy(window("event_time", "5 minutes"), "id")
.count()
.writeStream
.format("delta")
.option("checkpointLocation", "/checkpoint/path")
.outputMode("append")
.start()
)
# Volume grew, 200 partitions is no longer enough: scaling to 600
query.stop()
spark.conf.set("spark.sql.streaming.stateStore.partitions", "600")
query = (df
.withWatermark("event_time", "10 minutes")
.groupBy(window("event_time", "5 minutes"), "id")
.count()
.writeStream
.format("delta")
.option("checkpointLocation", "/checkpoint/path") # same checkpoint
.outputMode("append")
.start()
)The detail that matters here is that spark.sql.streaming.stateStore.partitions takes precedence over spark.sql.shuffle.partitions for that specific query, so there's no need to rewrite the aggregation logic or touch any other parameter, just restart with the checkpoint intact. In Lakeflow Pipelines, the same adjustment goes through spark_conf on the flow or table decorator, with no manual restart needed outside the pipeline's normal deploy cycle:
DP.append_flow(
target="session_table",
name="session_aggregation",
spark_conf={"spark.sql.streaming.stateStore.partitions": "600"}
)
def session_aggregation():
return (spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load(source_path)
.withWatermark("timestamp", "10 minutes")
.groupBy(window("timestamp", "5 minutes"), "id")
.count())The documented payoff: it's not just speed, it's the storage bill
The case Databricks itself published comes from Coveo, which reported a 40% reduction in Amazon S3 API call costs after being able to adjust state partitioning without migrating the entire checkpoint. That's an important clue about where the gain actually shows up: poorly sized state partitioning isn't just a processing-latency problem, it's also unnecessary read and write call volume against the object storage behind the checkpoint. Too few partitions concentrate too many calls on a single RocksDB instance; too many partitions spread coordination overhead without need. Both extremes cost money in a way that only really shows up on the storage bill, not on the latency dashboard.
How to know which number to pick (and how to confirm it worked)
Before triggering a repartition, it's worth measuring two numbers: the current state size per partition (visible in the stateOperators metric inside the StreamingQueryProgress event, the numRowsTotal field per operator) and the load distribution across partitions, since a key with much higher cardinality than the others concentrates disproportionate volume on a single RocksDB partition regardless of how many partitions exist in total. If the problem is total volume growing uniformly, increasing the partition count solves it. If the problem is one specific key dominating the volume (a hot key), increasing partitions alone won't help much, the bottleneck is in the key distribution, not the partition count.
After running the repartition, the same StreamingQueryProgress event that recorded the operation's duration in durationMs.controlBatch.REPARTITION goes back to reporting normal processing metrics on the next micro-batch. Comparing average micro-batch time before and after the change, together with storage read metrics if available in your observability setup, is the most direct way to confirm the adjustment had the expected effect, instead of assuming "more partitions is always better" and moving on without measuring.
What this doesn't solve
On-demand state repartitioning requires stopping the query to run, it isn't a live, real-time adjustment, so there's still a pause window proportional to the size of the accumulated state, larger state takes longer to redistribute. The feature also requires RocksDB as the state store provider, anyone still on the older in-memory provider (the legacy HDFS state store) needs to migrate before even considering this. And the feature is in Public Preview on Databricks Runtime 18, so before applying it to a critical production query, it's worth testing the repartition time against a checkpoint of comparable size in a non-production environment, the pause duration itself isn't documented as predictable in advance, we only know that "bigger state means longer."
Is it worth revisiting your configuration?
If your oldest stateful query has never had its partition count reviewed since creation, and data volume has grown since then (the common case), it's worth the exercise of measuring current state size and comparing it against the partition count inherited from day 1. The gain Coveo documented suggests that this kind of late adjustment, which used to require rebuilding the entire checkpoint (and for that reason almost never happened in practice), now has low enough operational cost to become a routine periodic review, not just an emergency measure.
References
- Databricks Docs, "On-demand state repartitioning for stateful streaming queries": https://docs.databricks.com/aws/en/structured-streaming/state-repartitioning
- Microsoft Learn, "On-demand state repartitioning for stateful streaming queries - Azure Databricks": https://learn.microsoft.com/en-us/azure/databricks/structured-streaming/state-repartitioning
- Databricks Blog, "Announcing On-Demand State Repartitioning for Apache Spark Structured Streaming on Databricks": https://www.databricks.com/blog/announcing-demand-state-repartitioning-apache-sparktm-structured-streaming-databricks