mithun-sudo commented on code in PR #18330:
URL: https://github.com/apache/iceberg/pull/18330#discussion_r4155254673


##########
flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicCommitter.java:
##########
@@ -284,6 +284,8 @@ private void replacePartitions(
 
     CommitSummary summary = new CommitSummary();
     summary.addAll(pendingResults);
+    Preconditions.checkState(
+        summary.deleteFilesCount() == 0, "Cannot overwrite partitions with 
delete files.");

Review Comment:
   Scenario: DynamicIcebergSink with overwrite enabled, while the Flink job 
still writes in upsert mode (update-by-key, e.g. CDC).
   
   Upsert makes the writer emit delete files plus data files. 
   
   Before this change, DynamicCommitter.replacePartitions() committed only the 
data files and ignored the deletes.
   
   IcebergCommitter already rejects that combo (Overwrite + upsert). This 
change makes DynamicCommitter fail the same way instead of succeeding silently. 
The goal here is to fail-fast.



-- 
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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to