Flink: reject delete files in dynamic sink overwrite mode - #18330
Open
mithun-sudo wants to merge 2 commits into
Open
mithun-sudo wants to merge 2 commits into
mithun-sudo wants to merge 2 commits into
Conversation
Align DynamicCommitter with IcebergCommitter by failing when replace partitions is requested but pending WriteResults contain delete files.
uros-b
reviewed
Oct 1, 2026
…rsions Port the DynamicCommitter guard and test to 2.2 and 1.20, and rename the test method per review feedback.
developer-rpai
approved these changes
Oct 2, 2026
developer-rpai
left a comment
There was a problem hiding this comment.
Parity fix I was glad to see — the dynamic committer silently dropping deletes on the overwrite path while the static committers reject them is exactly the kind of divergence that bites later. Two checks: is the error string byte-identical to the one in IcebergCommitter (one message for one condition)? And a layer question — the guard fires at commit time, after the delete files are already written, so a rejected commit leaves orphans for cleanup. If the static sink rejects at the same point, parity is the right call for now; worth noting if there's an earlier natural rejection point in the dynamic writer.
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.
Summary
Preconditions.checkStateguard inDynamicCommitter.replacePartitions()that already exists inIcebergCommitterandIcebergFilesCommitter.TestDynamicCommitter.testReplacePartitionsRejectsDeleteFilesto verify overwrite mode fails when delete files are present.Existing state
DynamicCommittersilently ignored delete files on the replace-partitions path while static sink committers reject that combination.Test plan
./gradlew :iceberg-flink:iceberg-flink-2.3:test --tests "org.apache.iceberg.flink.sink.dynamic.TestDynamicCommitter.testReplacePartitionsRejectsDeleteFiles"./gradlew :iceberg-flink:iceberg-flink-2.3:test --tests "org.apache.iceberg.flink.sink.dynamic.TestDynamicCommitter.testReplacePartitions"