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

Reply via email to