Hi Piotr,

Thanks for the proposal!

I have a couple of small questions:

1. Exactly-once sinks (let's say Kafka)
Will every sink subtask have exactly 1 transaction to commit per savepoint?
Will there be any transactions to abort by the streaming job? Or are the
transactions aborted on batch job restart?

2. Triggering checkpoint on a Task
Will that be done via RPC for non-source tasks? Or is it a part of the task
lifecycle?


Regards,
Roman


On Thu, Aug 13, 2026 at 8:12 AM Xuyang <[email protected]> wrote:

> 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