[
https://issues.apache.org/jira/browse/FLINK-40779?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18120202#comment-18120202
]
yongfu.gao commented on FLINK-40779:
------------------------------------
I'd like to pick this up.
The window is: {{enumerateSplits()}} advances the splitter cursor on a worker
thread, while the produced splits only reach {{unassigned}} in the handler on
the coordinator thread. A checkpoint taken in between records the advanced
state without those splits, so they are lost on restore.
I plan to fix it in {{JdbcSourceEnumerator}} so all {{SplitterEnumerator}}
implementations are covered and the {{@PublicEvolving}} interface stays
unchanged: capture splits and state together inside the callable, commit the
state in the handler, and use that committed state in {{snapshotState}}.
One thing to confirm: on the error path the committed state will lag, so a
restore may re-enumerate already-produced splits — at-least-once instead of
data loss. Is that acceptable?
A patch with a regression test is ready; I'll open a PR once the approach is
agreed.
> JdbcSourceEnumerator snapshots splitter state ahead of the splits it has
> handled
> --------------------------------------------------------------------------------
>
> Key: FLINK-40779
> URL: https://issues.apache.org/jira/browse/FLINK-40779
> Project: Flink
> Issue Type: Bug
> Components: Connectors / JDBC
> Affects Versions: jdbc-4.1.0
> Reporter: Martijn Visser
> Priority: Major
>
> {{JdbcSourceEnumerator}} runs {{SplitterEnumerator#enumerateSplits}} through
> {{callAsync}}, so the splitter advances its state on the worker thread while
> the splits only reach {{unassigned}} once the handler runs on the coordinator
> thread. A checkpoint taken in between stores the new splitter state without
> those splits, and they are lost on restore. {{SqlTemplateSplitEnumerator}} in
> continuous mode is affected today.
> {{snapshotState}} should use the splitter state that came back together with
> the last handled splits, captured inside the callable, instead of calling
> {{serializableState()}} directly.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)