[FLINK-40740][pipeline-connector/fluss] Prevent duplicate split initialization during discovery - #4555
Open
fxbing wants to merge 1 commit into
Open
[FLINK-40740][pipeline-connector/fluss] Prevent duplicate split initialization during discovery#4555fxbing wants to merge 1 commit into
fxbing wants to merge 1 commit into
Conversation
…alization during discovery Track initializing physical table paths before scheduling split creation and release each request in its coordinator callback without changing checkpoints. Add regressions for overlapping discovery, pending readers, partitions and recovery. Generated-by: Codex CLI 0.155.1
loserwang1024
self-requested a review
September 22, 2026 09:46
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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-40740.
Brief change log
Verifying this change
FlussSourceEnumeratorTest(16 tests) andFlussSourceEnumStateSerializerTest(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.Documentation
Was generative AI tooling used to co-author this PR?
Generated-by: Codex CLI 0.155.1