0
Accepted answer
Make sure cooperative rebalancing is on and you have enough cores to catch up. If one straggler partition lags, that's key skew, not partition count.
Structured Streaming job, 24 partitions, processing JSON events. During a rolling restart, lag climbs to around 40 minutes then recovers.
df.writeStream.format("delta").option("checkpointLocation", path).start()Should I tune maxOffsetsPerTrigger, add partitions, or just accept the redeploy lag?
Accepted answer
Make sure cooperative rebalancing is on and you have enough cores to catch up. If one straggler partition lags, that's key skew, not partition count.
Some lag during restart is expected while state stores recover. 40 minutes sounds high though. Check whether you have enough maxOffsetsPerTrigger headroom during steady state.
Sign in to reply.