fxbing opened a new pull request, #4555: URL: https://github.com/apache/flink-cdc/pull/4555
## What is the purpose of this pull request? Prevent periodic Fluss source discovery from initializing and assigning the same bucket splits more than once when a previous initialization callback has not completed. Reassigning an earliest-offset split can rewind an active reader and produce duplicate records. See FLINK-XXXXX. ## Brief change log - Track initializing physical table paths before scheduling asynchronous split creation, and release only the current batch in its coordinator callback. - Preserve initialization failures and the existing checkpoint format; unfinished initialization remains discoverable after recovery. - Add deterministic regressions for overlapping discovery, complete bucket assignment, independent partitions, pending readers, initialization failures, and recovery during initialization. ## Verifying this change - The overlapping-discovery regression fails against the original implementation: three buckets produce nine split assignments instead of three. - `FlussSourceEnumeratorTest` (16 tests) and `FlussSourceEnumStateSerializerTest` (1 test) pass on both Java 11 / Flink 1.20.3 and Java 17 / Flink 2.2.0, with no failures, errors, or skipped tests. - Spotless and Checkstyle pass. Full pipeline/CLI end-to-end tests were not run. ## Documentation - Does this pull request introduce a new feature? No. - If yes, how is the feature documented? Not applicable. ##### Was generative AI tooling used to co-author this PR? - [X] Yes — Codex CLI; Claude and GLM assisted with analysis, implementation, and review. Generated-by: Codex CLI 0.155.1 -- 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]
