Admaing opened a new pull request, #253:
URL: https://github.com/apache/flink-connector-jdbc/pull/253

   ## What is the purpose of the change
   
   `JdbcSourceEnumerator` enumerates splits through `context.callAsync(...)`, 
so the splitter cursor advances on a worker thread while the produced splits 
only reach `unassigned` in the handler on the coordinator thread. 
`snapshotState()` read the splitter state directly, so a checkpoint taken in 
that window recorded a state that had already moved past splits that were still 
in flight. Those splits were never re-enumerated after a restore.
   
   The fix captures the splits and the splitter state together inside the 
callable and commits the state from the handler; `snapshotState()` uses that 
committed state instead of calling `serializableState()`.
   
   The fix lives in the enumerator rather than in the individual 
`SplitterEnumerator` implementations, so every implementation is covered and 
the `@PublicEvolving` `SplitterEnumerator` interface stays unchanged. 
`SlideTimingSplitterEnumerator` is exposed the same way, since it advances 
`startMillis` during `enumerateSplits()`.
   
   ## Brief change log
   
     - `JdbcSourceEnumerator` keeps a `lastHandledSplitterState` field, 
committed by `onSplitsDiscovered` and read by `snapshotState()`
     - `enumerateSplitsWithState()` captures split and state on the same worker 
thread invocation and returns them as an `EnumerationResult`, so they can never 
be separated
     - `onSplitsDiscovered` commits the captured state unconditionally, 
including for empty enumerations, so a state that advanced the cursor is never 
left behind
     - `start()` seeds `lastHandledSplitterState` with the state as of start, 
which is what a checkpoint falls back to while the first batch is in flight
   
   ## Verifying this change
   
   This change added tests and can be verified as follows:
   
     - New `JdbcSourceEnumeratorCheckpointRaceTest` drives the async callable 
and its handler separately, takes a checkpoint in between, then restores and 
asserts the in-flight split is produced again instead of being dropped
     - The same test fails without the fix with `expected: 0 but was: 1`
     - The second case asserts the state is committed once the splits have been 
handled, and that the state of a not-yet-handled enumeration does not leak into 
the checkpoint
     - `JdbcSourceEnumeratorTest` still passes; 6/6 in 
`flink-connector-jdbc-core`
   
   Note for reviewers: `TestingSplitEnumeratorContext` runs the callable and 
its handler inside a single `Runnable`, so it cannot express this window. The 
new test overrides `callAsync` to keep the two apart — otherwise the race is 
not reproducible at all.
   
   ## Does this pull request potentially affect one of the following parts:
   
     - Dependencies (does it add or upgrade a dependency, including 
`flink.version` or a database driver version such as `mysql.version` / 
`postgres.version`): no
     - The public API, i.e., is any changed class annotated with 
`@Public(Evolving)` or `@Experimental`, or are the Table options or the JDBC 
dialect SPI (`JdbcDialect` / `JdbcFactory`) changed: no
     - Checkpointed state, its serializers, or exactly-once delivery (source 
splits, enumerator state, writer state, committables, XA transactions): yes
     - The per-record code paths (split reader, `JdbcOutputFormat`, dialect 
statement building; performance sensitive): no
   
   On the checkpointed state answer above: the serialized shape of 
`JdbcSourceEnumeratorState` is unchanged — only *which* splitter state is 
stored can differ, and only when a checkpoint lands in the window between an 
enumeration and its handler. That is precisely the bug being fixed: previously 
that checkpoint stored a state that had already passed splits which were not 
part of it. On the error path (the callable throws after partially advancing 
the state) the committed state now lags, so a restore may re-enumerate splits 
that were already produced. That is at-least-once rather than data loss, which 
seems the right direction for a source; happy to change it if you would rather 
keep the previous behaviour on failure.
   
   ## Documentation
   
     - Does this pull request introduce a new feature? no
     - If yes, how is the feature documented? not applicable
     - If the docs changed, are both `docs/content` and `docs/content.zh` 
updated? not applicable
   
   ---
   
   ##### Was generative AI tooling used to co-author this PR?
   
   - [X] Yes (DeepSeek Harness)
   
   Generated-by: DeepSeek Harness (deepseek-v4.1-flash)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to