BTW, could you please share the branch for the PoC as well as some details
about the benchmark (e.g., setup and results) ?
--
Best!
Xuyang
At 2026-08-12 14:22:22, "Xuyang" <[email protected]> wrote:
>Hi, Piotrek.
>Thanks a lot for putting this FLIP together. After modifying a streaming job,
>having to restart it statelessly often comes at a considerable time and
>resource cost for users. So big +1 from me.
>I've gone through the FLIP and it looks good to me overall. I just have a few
>questions I'd like to check with you:
>1. If the batch phase produces a savepoint in the standard format, my
>understanding is that the streaming phase could restore from it directly via
>`state.savepoints.dir`. If that's true, would we still need to set
>`execution.state-catch-up.mode: BATCH` in the streaming phase, or am I missing
>something here?
>2. Regarding sources, I'm wondering whether it might have an interface to
>indicate whether a source supports the bounded-to-unbounded transition. My
>thinking is that existing bounded sources in current connectors don't record
>their own state when they finish, so the subsequent streaming phase would have
>nothing to pick up from. Curious whether you've already considered how to
>handle this.
>3. I also have another rather rough idea around execution in the batch phase,
>and I'd be very interested in your thoughts on it. For very large full
>datasets, a single failure midway can invalidate hours of earlier batch
>computation. Even with fine-grained restarts in the batch phase down the road,
>a single full-scale run would still need to retain the blocking shuffle
>partitions for the entire backlog until they're consumed,so both TM disk
>pressure and the recomputation cost would scale with the full dataset. The
>idea would be to keep catch-up as a single job, but internally split the
>historical data into N shards(or N consecutive segments of data) that are
>consumed sequentially — producing a complete checkpoint as each finishes, and
>producing the save point handed over to the streaming phase when the last one
>completes. This feels somewhat in the spirit of what FLIP-309[1] describes,
>where checkpointing isn't disabled during backlog processing but simply
>happens at a lower frequency. That said, I may well be overlooking some
>complexity here, so please feel free to point out anything I've missed.
>Thanks again for the great work! Looking forward to your thoughts.
>
>
>[1]
>https://cwiki.apache.org/confluence/spaces/FLINK/pages/255069517/FLIP-309+Support+using+larger+checkpointing+interval+when+source+is+processing+backlog.
>
>
>
>--
>
> Best!
> Xuyang
>
>
>
>At 2026-08-07 21:53:16, "Piotr Nowojski" <[email protected]> wrote:
>>Hi all!
>>
>>I would like to start a discussion on state catch-up using BATCH execution
>>mode [1]
>>
>>FLIP-134 [2] introduced BATCH execution mode as a way to execute more
>>efficiently DataStream API jobs on bounded inputs. As even stated in the
>>documentation [3]:
>>
>>> One obvious outlier is when you want to use a bounded job to bootstrap
>>some job state that you
>>> then want to use in an unbounded job. For example, by running a bounded
>>job using
>>> STREAMING mode, taking a savepoint, and then restoring that savepoint on
>>an unbounded job.
>>> This is a very specific use case and one that might soon become obsolete
>>when we allow
>>> producing a savepoint as additional output of a BATCH execution job.
>>
>>producing a savepoint from BATCH execution job has always been planned as a
>>future extension. FLIP-605 [1] was created to implement this.
>>
>>The idea is that instead of waiting sometimes hours to process a long
>>backlog of records using STREAMING execution mode, users could submit their
>>jobs in BATCH execution mode, create a savepoint at the end, and then
>>recover their jobs from that savepoint and continue running them in
>>STREAMING. Regardless if those jobs are using DataStream API or Flink
>>Streaming SQL/Table API.
>>
>>For more information please take a look into the FLIP document itself [1].
>>
>>Disclosure, in Confluent/IBM we have a working early access version of this
>>feature, hence we already have a pretty good understanding of what needs to
>>be done to make it work and what are the expected results.
>>
>>Best,
>>Piotrek
>>
>>[1]
>>https://cwiki.apache.org/confluence/spaces/FLINK/pages/446071315/FLIP-605+State+catch-up+of+streaming+jobs+state+using+BATCH+execution+mode
>>[2]
>>https://cwiki.apache.org/confluence/spaces/FLINK/pages/158871522/FLIP-134+Batch+execution+for+the+DataStream+API
>>[3]
>>https://nightlies.apache.org/flink/flink-docs-master/docs/dev/datastream/execution_mode/#when-canshould-i-use-batch-execution-mode