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]