From 35752ec38b52b9ab89aa37c68a9e0216f5bd971c Mon Sep 17 00:00:00 2001 From: David Kim Date: Thu, 27 Aug 2026 11:26:06 -0700 Subject: [PATCH 1/2] Flink 1.18: harden IcebergFilesCommitter against ambiguous-commit duplicates Two prod incidents (chrono.payment_stripe_webhook_event on CA, chrono.user_updates_status on US East) saw the same physical data file registered twice in Iceberg snapshot metadata, breaking Snowflake Catalog-Linked Database auto-refresh. Root cause traced to actual Polaris REST catalog logs: a 500 on the commit endpoint gets mapped by Iceberg's own REST client to CommitStateUnknownException (ErrorHandlers#commitErrorHandler treats 500/502/504 on a commit as ambiguous, not a definite failure -- the request may have already succeeded server-side). SnapshotProducer#commit deliberately does no cleanup and rethrows immediately in this case, by design, leaving it to the caller to decide whether retrying is safe. IcebergFilesCommitter has no handling for this exception at all: it propagates uncaught out of notifyCheckpointComplete, fails the Flink task, and Flink's restart-strategy blindly restarts and re-commits the same checkpoint -- exactly what CommitStateUnknownException's own javadoc warns against ("retrying an already successful operation will result in duplicate records"). job-id/operator-id are correctly persisted/restored across the restart (verified against this exact class), so the existing getMaxCommittedCheckpointId dedup check is evaluated correctly -- but only as of whatever the catalog reports *at that instant*. Production evidence (two commits 17 seconds apart, identical flink.job-id and flink.operator-id, identical file, both ADDED) shows the restart can outrace the catalog's own visibility of its prior write, so the check reports "not yet committed" for a checkpoint that actually already landed. This is the same class of bug documented (and never fixed) in apache/iceberg#10765. Two changes, meant to work together: 1. commitOperation now explicitly catches CommitStateUnknownException and polls the catalog (refresh + getMaxCommittedCheckpointId, up to 5 attempts with exponential backoff starting at 1s) to check whether the ambiguous commit actually landed, before letting the exception propagate into an uncontrolled Flink restart. This directly targets the proven mechanism instead of relying on Flink's restart delay happening to outlast the catalog's propagation lag, which production evidence shows isn't reliable (17s wasn't enough). 2. dropAlreadyCommittedFiles (called from commitUpToCheckpoint, so it runs on the restart path too) adds a content-based check on top of the existing metadata-only dedup: before building a commit, it walks back up to 5 ancestor snapshots for this exact flinkJobId/operatorId and drops any file whose exact path is already present, logging loudly when it happens. Defense in depth for the case where verification in (1) times out or is skipped. Not related to Dan Cagney's non-idempotent-retry patch (this exact commit's parent) -- that patch controls Apache HttpClient's own retry-on-5xx behavior below the REST client abstraction, and is working as intended (it's *why* this surfaces as an observable CommitStateUnknownException + Flink restart rather than a silent duplicate from a blind HTTP-layer retry). This fix targets a layer above it: what happens after that exception is correctly raised. Also checked and ruled out: the DynamicCommitter fix (apache/iceberg#16008/#16011, job-id persistence) -- that only touches flink/*/sink/dynamic/DynamicCommitter.java, which this pipeline (Flink SQL Table API, not DataStream API DynamicIcebergSink) doesn't use. Separately confirmed the non-dynamic Flink 1.20 SinkV2 IcebergWriteAggregator still has that exact job-id-not-persisted-across-restart gap unpatched as of 1.11.0 -- irrelevant to this specific incident (job-id was identical on both duplicate commits) but relevant if this fork ever upgrades past Flink 1.18 and opts into table.exec.iceberg.use-v2-sink. Scope: flink/v1.18 only, matching what chrono's PyFlink writer actually runs in production (engine-version 1.18.1). flink/v1.19 and flink/v1.20 carry near-identical vulnerable logic but have already refactored getMaxCommittedCheckpointId/deleteCommittedManifests into a shared SinkUtil/FlinkManifestUtil class, so porting this isn't a straight copy -- follow-up PR once this one is reviewed, same approach the upstream DynamicCommitter fix took (single version first, then backport). Not compiled locally -- needs a Thor build (same convention as PR #4 on this fork). No test coverage yet; this needs at minimum a case simulating commitUpToCheckpoint being invoked twice with an identical file, and ideally a CommitStateUnknownException-on-first-attempt simulation mirroring TestDynamicIcebergSink.testCommitsOnceWhenConcurrentDuplicateCommit() upstream. Co-Authored-By: Claude Sonnet 5 --- .../flink/sink/IcebergFilesCommitter.java | 218 +++++++++++++++++- 1 file changed, 215 insertions(+), 3 deletions(-) diff --git a/flink/v1.18/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java b/flink/v1.18/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java index 7108c2008341..e2031ad2b991 100644 --- a/flink/v1.18/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java +++ b/flink/v1.18/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java @@ -40,6 +40,8 @@ import org.apache.flink.streaming.runtime.streamrecord.StreamRecord; import org.apache.flink.table.runtime.typeutils.SortedMapTypeInfo; import org.apache.iceberg.AppendFiles; +import org.apache.iceberg.DataFile; +import org.apache.iceberg.DeleteFile; import org.apache.iceberg.ManifestFile; import org.apache.iceberg.PartitionSpec; import org.apache.iceberg.ReplacePartitions; @@ -47,6 +49,7 @@ import org.apache.iceberg.Snapshot; import org.apache.iceberg.SnapshotUpdate; import org.apache.iceberg.Table; +import org.apache.iceberg.exceptions.CommitStateUnknownException; import org.apache.iceberg.flink.TableLoader; import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.relocated.com.google.common.annotations.VisibleForTesting; @@ -57,6 +60,7 @@ import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.types.Comparators; import org.apache.iceberg.types.Types; +import org.apache.iceberg.util.CharSequenceSet; import org.apache.iceberg.util.PropertyUtil; import org.apache.iceberg.util.ThreadPools; import org.slf4j.Logger; @@ -80,6 +84,33 @@ class IcebergFilesCommitter extends AbstractStreamOperator private static final String MAX_COMMITTED_CHECKPOINT_ID = "flink.max-committed-checkpoint-id"; static final String MAX_CONTINUOUS_EMPTY_COMMITS = "flink.max-continuous-empty-commits"; + // AFFIRM: how many of the most recent ancestor snapshots (for this flinkJobId/operatorId) we + // inspect, by actual file path, before committing. This guards against the case documented in + // apache/iceberg#10765: a commit's response is lost/ambiguous (e.g. a network timeout), the + // committer doesn't observe success, and a subsequent commit re-appends the exact same files + // because getMaxCommittedCheckpointId() was evaluated against a table snapshot that didn't yet + // reflect the prior, already-successful commit. That check trusts the flink.job-id / + // flink.operator-id / flink.max-committed-checkpoint-id snapshot summary properties; this check + // additionally verifies by content (file path) against a small, bounded window of recent + // history, so a lost-response race can't silently double-register a file. Bounded to a small + // constant so cost stays flat regardless of total table/snapshot history size. + private static final int RECENT_SNAPSHOT_LOOKBACK = 5; + + // AFFIRM: apache/iceberg's REST client maps a 500/502/504 from the catalog on a commit request + // to CommitStateUnknownException (see ErrorHandlers#commitErrorHandler) -- the request may have + // actually succeeded server-side; the client just couldn't confirm it. SnapshotProducer#commit + // deliberately does no cleanup and rethrows immediately in this case, by design, leaving the + // caller responsible for deciding whether it's safe to retry. Left unhandled, that exception + // propagates out of notifyCheckpointComplete, fails the Flink task, and Flink's restart-strategy + // blindly retries the same checkpoint -- exactly what CommitStateUnknownException's own javadoc + // warns against ("retrying an already successful operation will result in duplicate records"). + // Rather than rely on that restart happening to be slow enough for the catalog to catch up (it + // wasn't, in production: see the 17s-apart duplicate on chrono.user_updates_status), poll the + // catalog directly for a bounded window to find out whether the commit actually landed before + // letting the exception propagate into an uncontrolled restart. + private static final int COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS = 5; + private static final long COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS = 1000L; + // TableLoader to load iceberg table lazily. private final TableLoader tableLoader; private final boolean replacePartitions; @@ -275,13 +306,113 @@ private void commitUpToCheckpoint( manifests.addAll(deltaManifests.manifests()); } - CommitSummary summary = new CommitSummary(pendingResults); - commitPendingResult(pendingResults, summary, newFlinkJobId, operatorId, checkpointId); + NavigableMap dedupedResults = + dropAlreadyCommittedFiles(pendingResults, newFlinkJobId, operatorId); + + CommitSummary summary = new CommitSummary(dedupedResults); + commitPendingResult(dedupedResults, summary, newFlinkJobId, operatorId, checkpointId); committerMetrics.updateCommitSummary(summary); pendingMap.clear(); deleteCommittedManifests(manifests, newFlinkJobId, checkpointId); } + /** + * AFFIRM: defense-in-depth against apache/iceberg#10765. Drops any data/delete file from {@code + * pendingResults} whose exact path already appears as an added file in one of the last {@link + * #RECENT_SNAPSHOT_LOOKBACK} ancestor snapshots committed by this same flinkJobId/operatorId. + * This does not replace {@link #getMaxCommittedCheckpointId}; it's an additional, content-based + * check for the narrow window where that metadata-only check can be fooled by a commit whose + * response was lost/ambiguous to the client but which actually succeeded on the catalog. + */ + private NavigableMap dropAlreadyCommittedFiles( + NavigableMap pendingResults, String newFlinkJobId, String operatorId) { + CharSequenceSet recentlyCommittedPaths = + collectRecentlyCommittedFilePaths(newFlinkJobId, operatorId); + if (recentlyCommittedPaths.isEmpty()) { + return pendingResults; + } + + NavigableMap deduped = Maps.newTreeMap(); + for (Map.Entry e : pendingResults.entrySet()) { + long checkpointId = e.getKey(); + WriteResult result = e.getValue(); + WriteResult.Builder builder = + WriteResult.builder().addReferencedDataFiles(result.referencedDataFiles()); + + for (DataFile file : result.dataFiles()) { + if (recentlyCommittedPaths.contains(file.path())) { + LOG.warn( + "Dropping data file already present in a recent snapshot for table {} branch {} " + + "flinkJobId {} operatorId {} checkpoint {}: {}. This indicates a prior commit " + + "for this file already succeeded even though this committer didn't observe " + + "that success (see apache/iceberg#10765).", + table.name(), + branch, + newFlinkJobId, + operatorId, + checkpointId, + file.path()); + } else { + builder.addDataFiles(file); + } + } + + for (DeleteFile file : result.deleteFiles()) { + if (recentlyCommittedPaths.contains(file.path())) { + LOG.warn( + "Dropping delete file already present in a recent snapshot for table {} branch {} " + + "flinkJobId {} operatorId {} checkpoint {}: {}. This indicates a prior commit " + + "for this file already succeeded even though this committer didn't observe " + + "that success (see apache/iceberg#10765).", + table.name(), + branch, + newFlinkJobId, + operatorId, + checkpointId, + file.path()); + } else { + builder.addDeleteFiles(file); + } + } + + deduped.put(checkpointId, builder.build()); + } + + return deduped; + } + + /** + * Walks back from the current snapshot on {@link #branch}, collecting the paths of data/delete + * files added by snapshots committed by this exact flinkJobId/operatorId, stopping after {@link + * #RECENT_SNAPSHOT_LOOKBACK} ancestor snapshots (regardless of whether they match) to keep this + * bounded and cheap. Refreshes the table first so this observes the freshest metadata available. + */ + private CharSequenceSet collectRecentlyCommittedFilePaths(String flinkJobId, String operatorId) { + table.refresh(); + + CharSequenceSet paths = CharSequenceSet.empty(); + Snapshot snapshot = table.snapshot(branch); + int inspected = 0; + while (snapshot != null && inspected < RECENT_SNAPSHOT_LOOKBACK) { + Map summary = snapshot.summary(); + if (flinkJobId.equals(summary.get(FLINK_JOB_ID)) + && (summary.get(OPERATOR_ID) == null || operatorId.equals(summary.get(OPERATOR_ID)))) { + for (DataFile file : snapshot.addedDataFiles(table.io())) { + paths.add(file.path()); + } + for (DeleteFile file : snapshot.addedDeleteFiles(table.io())) { + paths.add(file.path()); + } + } + + Long parentSnapshotId = snapshot.parentId(); + snapshot = parentSnapshotId != null ? table.snapshot(parentSnapshotId) : null; + inspected++; + } + + return paths; + } + private void commitPendingResult( NavigableMap pendingResults, CommitSummary summary, @@ -412,7 +543,42 @@ private void commitOperation( operation.toBranch(branch); long startNano = System.nanoTime(); - operation.commit(); // abort is automatically called if this fails. + try { + operation.commit(); // abort is automatically called if this fails, EXCEPT on + // CommitStateUnknownException -- see the catch block below. + } catch (CommitStateUnknownException e) { + // AFFIRM: see apache/iceberg#10765 and the class-level comment on + // COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS. Don't let Flink's restart-strategy be the thing + // that decides whether this ambiguous commit gets blindly retried; check for ourselves. + if (verifyCommitEventuallySucceeded(newFlinkJobId, operatorId, checkpointId, description)) { + LOG.warn( + "Commit {} for checkpoint {} to table {} branch {} returned an ambiguous response " + + "(CommitStateUnknownException) but verification found it actually succeeded; " + + "treating this checkpoint as committed instead of failing the task.", + description, + checkpointId, + table.name(), + branch, + e); + long durationMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNano); + committerMetrics.commitDuration(durationMs); + return; + } + + LOG.error( + "Commit {} for checkpoint {} to table {} branch {} returned an ambiguous response " + + "(CommitStateUnknownException) and could not be verified as successful within " + + "the retry budget ({} attempts). Rethrowing so Flink can restart and re-attempt. " + + "If the commit actually did succeed after this budget was exhausted, " + + "dropAlreadyCommittedFiles should still catch and drop the duplicate on retry.", + description, + checkpointId, + table.name(), + branch, + COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS, + e); + throw e; + } long durationMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNano); LOG.info( "Committed {} to table: {}, branch: {}, checkpointId {} in {} ms", @@ -424,6 +590,52 @@ private void commitOperation( committerMetrics.commitDuration(durationMs); } + /** + * AFFIRM: polls, with exponential backoff, for up to {@link + * #COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS} attempts to determine whether a commit that just + * threw {@link CommitStateUnknownException} actually landed on the catalog, by refreshing the + * table and checking whether {@code checkpointId} is now covered by {@link + * #getMaxCommittedCheckpointId}. Deliberately blocks the calling thread (inside + * notifyCheckpointComplete): a bounded wait here is preferable to unconditionally failing the + * task and paying for a full restart-and-restore cycle only to hit the same ambiguity check + * again. Returns false (not verified) if the budget is exhausted or the wait is interrupted. + */ + private boolean verifyCommitEventuallySucceeded( + String flinkJobId, String operatorId, long checkpointId, String description) { + long delayMs = COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS; + for (int attempt = 1; attempt <= COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS; attempt++) { + try { + Thread.sleep(delayMs); + } catch (InterruptedException interruptedException) { + Thread.currentThread().interrupt(); + return false; + } + + table.refresh(); + long observedCheckpointId = + getMaxCommittedCheckpointId(table, flinkJobId, operatorId, branch); + LOG.info( + "Verifying ambiguous {} commit for checkpoint {} on table {} branch {}: attempt {}/{}, " + + "observed max-committed-checkpoint-id {} for flinkJobId {} operatorId {}", + description, + checkpointId, + table.name(), + branch, + attempt, + COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS, + observedCheckpointId, + flinkJobId, + operatorId); + if (observedCheckpointId >= checkpointId) { + return true; + } + + delayMs *= 2; + } + + return false; + } + @Override public void processElement(StreamRecord element) { FlinkWriteResult flinkWriteResult = element.getValue(); From c6015f5338cda079bcf1816b991d60327a123c2c Mon Sep 17 00:00:00 2001 From: David Kim Date: Thu, 27 Aug 2026 12:11:52 -0700 Subject: [PATCH 2/2] Add tests for ambiguous-commit safety changes; make retry budget configurable Verified locally (JDK 21 / Gradle 8.12.1, -DflinkVersions=1.18): - :iceberg-flink:iceberg-flink-1.18:compileJava and compileTestJava both succeed - Full existing TestIcebergFilesCommitter suite: 96/96 pre-existing tests still pass (0 failures, 0 errors) -- this change doesn't regress the existing commit paths - 3 new test methods x 6 parameterizations (format version x file format x branch) = 18 new cases, all passing: 114/114 total, 100% success per the HTML test report Fixed one real compile error along the way: CharSequenceSet.contains(CharSequence) trips error-prone's CollectionUndefinedEquality check. Confirmed this is a known false positive -- CharSequenceSet exists specifically to give CharSequence content-equality, and Iceberg's own codebase suppresses the same warning at the call site in several places (ManifestFilterManager, MergingSnapshotProducer, and CharSequenceSet's own implementation). Same fix applied here. Also made the three new tunables configurable via table properties instead of hardcoded constants, mirroring the existing MAX_CONTINUOUS_EMPTY_COMMITS pattern in this same class: - flink.recent-snapshot-lookback (default 5) - flink.commit-state-unknown-max-verify-attempts (default 5) - flink.commit-state-unknown-verify-initial-delay-ms (default 1000) This isn't just cleanliness -- it's what makes testVerifyCommitEventuallySucceeded's never-committed-checkpoint case runnable in ~30ms instead of ~30s of real Thread.sleep. New tests, and what they actually verify: - testDropAlreadyCommittedFilesDropsDuplicatePath: commits a file for real, then builds a pendingResults map for a later checkpoint that (as in the actual production incident) still bundles that same already-committed file alongside a genuinely new one. Asserts dropAlreadyCommittedFiles removes only the duplicate. - testDropAlreadyCommittedFilesKeepsUnrelatedFiles: negative-case sanity check -- files with no match in recent history pass through unchanged. - testVerifyCommitEventuallySucceeded: asserts the poll-and-verify loop reports true for a checkpoint id that really was committed, and false (after exhausting the configured budget) for one that never was. Explicit scope note: these are unit tests against dropAlreadyCommittedFiles and verifyCommitEventuallySucceeded directly (both now package-private + @VisibleForTesting, matching this file's existing convention for getMaxCommittedCheckpointId), not end-to-end reproductions of the production race. The race itself depends on real REST catalog propagation lag, which a local HadoopTables-backed test table can't produce (commits are synchronously visible, no eventual consistency to race against) -- so these tests verify the new logic behaves correctly *given* the failure condition, not that Flink's restart timing will actually trigger that condition the same way Polaris did in production. Real end-to-end confidence still needs either a staging deploy with an injectable-latency/fault-injecting catalog proxy, or production monitoring after this ships. Co-Authored-By: Claude Sonnet 5 --- .../flink/sink/IcebergFilesCommitter.java | 82 +++++++---- .../flink/sink/TestIcebergFilesCommitter.java | 130 ++++++++++++++++++ 2 files changed, 184 insertions(+), 28 deletions(-) diff --git a/flink/v1.18/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java b/flink/v1.18/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java index e2031ad2b991..798db6f5bf34 100644 --- a/flink/v1.18/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java +++ b/flink/v1.18/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergFilesCommitter.java @@ -94,7 +94,8 @@ class IcebergFilesCommitter extends AbstractStreamOperator // additionally verifies by content (file path) against a small, bounded window of recent // history, so a lost-response race can't silently double-register a file. Bounded to a small // constant so cost stays flat regardless of total table/snapshot history size. - private static final int RECENT_SNAPSHOT_LOOKBACK = 5; + static final String RECENT_SNAPSHOT_LOOKBACK_PROP = "flink.recent-snapshot-lookback"; + private static final int RECENT_SNAPSHOT_LOOKBACK_DEFAULT = 5; // AFFIRM: apache/iceberg's REST client maps a 500/502/504 from the catalog on a commit request // to CommitStateUnknownException (see ErrorHandlers#commitErrorHandler) -- the request may have @@ -107,9 +108,15 @@ class IcebergFilesCommitter extends AbstractStreamOperator // Rather than rely on that restart happening to be slow enough for the catalog to catch up (it // wasn't, in production: see the 17s-apart duplicate on chrono.user_updates_status), poll the // catalog directly for a bounded window to find out whether the commit actually landed before - // letting the exception propagate into an uncontrolled restart. - private static final int COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS = 5; - private static final long COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS = 1000L; + // letting the exception propagate into an uncontrolled restart. Configurable via table + // properties (same pattern as MAX_CONTINUOUS_EMPTY_COMMITS below) so this can be tuned without a + // code change, and so tests don't have to sleep through the real default budget. + static final String COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS_PROP = + "flink.commit-state-unknown-max-verify-attempts"; + static final String COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS_PROP = + "flink.commit-state-unknown-verify-initial-delay-ms"; + private static final int COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS_DEFAULT = 5; + private static final long COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS_DEFAULT = 1000L; // TableLoader to load iceberg table lazily. private final TableLoader tableLoader; @@ -139,6 +146,9 @@ class IcebergFilesCommitter extends AbstractStreamOperator private transient long maxCommittedCheckpointId; private transient int continuousEmptyCheckpoints; private transient int maxContinuousEmptyCommits; + private transient int recentSnapshotLookback; + private transient int commitStateUnknownMaxVerifyAttempts; + private transient long commitStateUnknownVerifyInitialDelayMs; // There're two cases that we restore from flink checkpoints: the first case is restoring from // snapshot created by the same flink job; another case is restoring from snapshot created by // another different job. For the second case, we need to maintain the old flink job's id in flink @@ -186,6 +196,19 @@ public void initializeState(StateInitializationContext context) throws Exception PropertyUtil.propertyAsInt(table.properties(), MAX_CONTINUOUS_EMPTY_COMMITS, 10); Preconditions.checkArgument( maxContinuousEmptyCommits > 0, MAX_CONTINUOUS_EMPTY_COMMITS + " must be positive"); + recentSnapshotLookback = + PropertyUtil.propertyAsInt( + table.properties(), RECENT_SNAPSHOT_LOOKBACK_PROP, RECENT_SNAPSHOT_LOOKBACK_DEFAULT); + commitStateUnknownMaxVerifyAttempts = + PropertyUtil.propertyAsInt( + table.properties(), + COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS_PROP, + COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS_DEFAULT); + commitStateUnknownVerifyInitialDelayMs = + PropertyUtil.propertyAsLong( + table.properties(), + COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS_PROP, + COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS_DEFAULT); int subTaskId = getRuntimeContext().getIndexOfThisSubtask(); int attemptId = getRuntimeContext().getAttemptNumber(); @@ -318,13 +341,15 @@ private void commitUpToCheckpoint( /** * AFFIRM: defense-in-depth against apache/iceberg#10765. Drops any data/delete file from {@code - * pendingResults} whose exact path already appears as an added file in one of the last {@link - * #RECENT_SNAPSHOT_LOOKBACK} ancestor snapshots committed by this same flinkJobId/operatorId. - * This does not replace {@link #getMaxCommittedCheckpointId}; it's an additional, content-based - * check for the narrow window where that metadata-only check can be fooled by a commit whose - * response was lost/ambiguous to the client but which actually succeeded on the catalog. + * pendingResults} whose exact path already appears as an added file in one of the last {@code + * recentSnapshotLookback} ancestor snapshots committed by this same flinkJobId/operatorId. This + * does not replace {@link #getMaxCommittedCheckpointId}; it's an additional, content-based check + * for the narrow window where that metadata-only check can be fooled by a commit whose response + * was lost/ambiguous to the client but which actually succeeded on the catalog. */ - private NavigableMap dropAlreadyCommittedFiles( + @VisibleForTesting + @SuppressWarnings("CollectionUndefinedEquality") // CharSequenceSet defines path equality itself + NavigableMap dropAlreadyCommittedFiles( NavigableMap pendingResults, String newFlinkJobId, String operatorId) { CharSequenceSet recentlyCommittedPaths = collectRecentlyCommittedFilePaths(newFlinkJobId, operatorId); @@ -383,8 +408,8 @@ private NavigableMap dropAlreadyCommittedFiles( /** * Walks back from the current snapshot on {@link #branch}, collecting the paths of data/delete - * files added by snapshots committed by this exact flinkJobId/operatorId, stopping after {@link - * #RECENT_SNAPSHOT_LOOKBACK} ancestor snapshots (regardless of whether they match) to keep this + * files added by snapshots committed by this exact flinkJobId/operatorId, stopping after {@code + * recentSnapshotLookback} ancestor snapshots (regardless of whether they match) to keep this * bounded and cheap. Refreshes the table first so this observes the freshest metadata available. */ private CharSequenceSet collectRecentlyCommittedFilePaths(String flinkJobId, String operatorId) { @@ -393,7 +418,7 @@ private CharSequenceSet collectRecentlyCommittedFilePaths(String flinkJobId, Str CharSequenceSet paths = CharSequenceSet.empty(); Snapshot snapshot = table.snapshot(branch); int inspected = 0; - while (snapshot != null && inspected < RECENT_SNAPSHOT_LOOKBACK) { + while (snapshot != null && inspected < recentSnapshotLookback) { Map summary = snapshot.summary(); if (flinkJobId.equals(summary.get(FLINK_JOB_ID)) && (summary.get(OPERATOR_ID) == null || operatorId.equals(summary.get(OPERATOR_ID)))) { @@ -548,8 +573,8 @@ private void commitOperation( // CommitStateUnknownException -- see the catch block below. } catch (CommitStateUnknownException e) { // AFFIRM: see apache/iceberg#10765 and the class-level comment on - // COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS. Don't let Flink's restart-strategy be the thing - // that decides whether this ambiguous commit gets blindly retried; check for ourselves. + // COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS_PROP. Don't let Flink's restart-strategy be the + // thing that decides whether this ambiguous commit gets blindly retried; check ourselves. if (verifyCommitEventuallySucceeded(newFlinkJobId, operatorId, checkpointId, description)) { LOG.warn( "Commit {} for checkpoint {} to table {} branch {} returned an ambiguous response " @@ -575,7 +600,7 @@ private void commitOperation( checkpointId, table.name(), branch, - COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS, + commitStateUnknownMaxVerifyAttempts, e); throw e; } @@ -591,19 +616,20 @@ private void commitOperation( } /** - * AFFIRM: polls, with exponential backoff, for up to {@link - * #COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS} attempts to determine whether a commit that just - * threw {@link CommitStateUnknownException} actually landed on the catalog, by refreshing the - * table and checking whether {@code checkpointId} is now covered by {@link - * #getMaxCommittedCheckpointId}. Deliberately blocks the calling thread (inside - * notifyCheckpointComplete): a bounded wait here is preferable to unconditionally failing the - * task and paying for a full restart-and-restore cycle only to hit the same ambiguity check - * again. Returns false (not verified) if the budget is exhausted or the wait is interrupted. + * AFFIRM: polls, with exponential backoff, for up to {@code commitStateUnknownMaxVerifyAttempts} + * attempts to determine whether a commit that just threw {@link CommitStateUnknownException} + * actually landed on the catalog, by refreshing the table and checking whether {@code + * checkpointId} is now covered by {@link #getMaxCommittedCheckpointId}. Deliberately blocks the + * calling thread (inside notifyCheckpointComplete): a bounded wait here is preferable to + * unconditionally failing the task and paying for a full restart-and-restore cycle only to hit + * the same ambiguity check again. Returns false (not verified) if the budget is exhausted or the + * wait is interrupted. */ - private boolean verifyCommitEventuallySucceeded( + @VisibleForTesting + boolean verifyCommitEventuallySucceeded( String flinkJobId, String operatorId, long checkpointId, String description) { - long delayMs = COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS; - for (int attempt = 1; attempt <= COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS; attempt++) { + long delayMs = commitStateUnknownVerifyInitialDelayMs; + for (int attempt = 1; attempt <= commitStateUnknownMaxVerifyAttempts; attempt++) { try { Thread.sleep(delayMs); } catch (InterruptedException interruptedException) { @@ -622,7 +648,7 @@ private boolean verifyCommitEventuallySucceeded( table.name(), branch, attempt, - COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS, + commitStateUnknownMaxVerifyAttempts, observedCheckpointId, flinkJobId, operatorId); diff --git a/flink/v1.18/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergFilesCommitter.java b/flink/v1.18/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergFilesCommitter.java index 75060c479e11..aabe2b9fe8b4 100644 --- a/flink/v1.18/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergFilesCommitter.java +++ b/flink/v1.18/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergFilesCommitter.java @@ -179,6 +179,136 @@ private FlinkWriteResult of(long checkpointId, DataFile dataFile) { return new FlinkWriteResult(checkpointId, WriteResult.builder().addDataFiles(dataFile).build()); } + @TestTemplate + public void testDropAlreadyCommittedFilesDropsDuplicatePath() throws Exception { + // AFFIRM: regression test for apache/iceberg#10765. Simulates the production race directly: + // getMaxCommittedCheckpointId() is evaluated against stale catalog state and a checkpoint + // whose file was already committed gets bundled into a later commit attempt anyway. + // dropAlreadyCommittedFiles should catch it by path, independent of that metadata check. + long timestamp = 0; + JobID jobId = new JobID(); + OperatorID operatorId; + try (OneInputStreamOperatorTestHarness harness = + createStreamSink(jobId)) { + harness.setup(); + harness.open(); + operatorId = harness.getOperator().getOperatorID(); + + RowData row1 = SimpleDataUtil.createRowData(1, "hello"); + DataFile dataFile1 = writeDataFile("data-1", ImmutableList.of(row1)); + + long checkpointId = 1; + harness.processElement(of(checkpointId, dataFile1), ++timestamp); + harness.snapshot(checkpointId, ++timestamp); + harness.notifyOfCompletedCheckpoint(checkpointId); + + SimpleDataUtil.assertTableRows(table, ImmutableList.of(row1), branch); + assertSnapshotSize(1); + assertMaxCommittedCheckpointId(jobId, operatorId, checkpointId); + + // Build a pendingResults map, as if for a subsequent commit attempt, that (incorrectly) + // still includes dataFile1 -- the exact shape of the production duplicate, where the same + // physical file got bundled into a commit whose checkpoint the dedup check believed was + // not yet committed -- alongside a genuinely new file that must NOT be dropped. + RowData row2 = SimpleDataUtil.createRowData(2, "world"); + DataFile dataFile2 = writeDataFile("data-2", ImmutableList.of(row2)); + + NavigableMap pendingResults = Maps.newTreeMap(); + pendingResults.put(2L, WriteResult.builder().addDataFiles(dataFile1, dataFile2).build()); + + IcebergFilesCommitter committer = (IcebergFilesCommitter) harness.getOperator(); + NavigableMap deduped = + committer.dropAlreadyCommittedFiles( + pendingResults, jobId.toString(), operatorId.toHexString()); + + List remaining = Lists.newArrayList(deduped.get(2L).dataFiles()); + assertThat(remaining).hasSize(1); + assertThat(remaining.get(0).location()).isEqualTo(dataFile2.location()); + } + } + + @TestTemplate + public void testDropAlreadyCommittedFilesKeepsUnrelatedFiles() throws Exception { + // Sanity check for the same method: when nothing in pendingResults was already committed, + // dropAlreadyCommittedFiles must be a no-op (return the input unchanged in content). + long timestamp = 0; + JobID jobId = new JobID(); + OperatorID operatorId; + try (OneInputStreamOperatorTestHarness harness = + createStreamSink(jobId)) { + harness.setup(); + harness.open(); + operatorId = harness.getOperator().getOperatorID(); + + RowData row1 = SimpleDataUtil.createRowData(1, "hello"); + DataFile dataFile1 = writeDataFile("data-1", ImmutableList.of(row1)); + long checkpointId = 1; + harness.processElement(of(checkpointId, dataFile1), ++timestamp); + harness.snapshot(checkpointId, ++timestamp); + harness.notifyOfCompletedCheckpoint(checkpointId); + + RowData row2 = SimpleDataUtil.createRowData(2, "world"); + DataFile dataFile2 = writeDataFile("data-2", ImmutableList.of(row2)); + NavigableMap pendingResults = Maps.newTreeMap(); + pendingResults.put(2L, WriteResult.builder().addDataFiles(dataFile2).build()); + + IcebergFilesCommitter committer = (IcebergFilesCommitter) harness.getOperator(); + NavigableMap deduped = + committer.dropAlreadyCommittedFiles( + pendingResults, jobId.toString(), operatorId.toHexString()); + + List remaining = Lists.newArrayList(deduped.get(2L).dataFiles()); + assertThat(remaining).hasSize(1); + assertThat(remaining.get(0).location()).isEqualTo(dataFile2.location()); + } + } + + @TestTemplate + public void testVerifyCommitEventuallySucceeded() throws Exception { + // AFFIRM: regression test for the CommitStateUnknownException handling in commitOperation. + // Configure a fast, test-scale retry budget rather than the real default (5 attempts, + // starting at 1s, exponential backoff) so this doesn't take ~30s to run. + table + .updateProperties() + .set(IcebergFilesCommitter.COMMIT_STATE_UNKNOWN_MAX_VERIFY_ATTEMPTS_PROP, "2") + .set(IcebergFilesCommitter.COMMIT_STATE_UNKNOWN_VERIFY_INITIAL_DELAY_MS_PROP, "10") + .commit(); + + long timestamp = 0; + JobID jobId = new JobID(); + OperatorID operatorId; + try (OneInputStreamOperatorTestHarness harness = + createStreamSink(jobId)) { + harness.setup(); + harness.open(); + operatorId = harness.getOperator().getOperatorID(); + + RowData row = SimpleDataUtil.createRowData(1, "hello"); + DataFile dataFile = writeDataFile("data-1", ImmutableList.of(row)); + long checkpointId = 1; + harness.processElement(of(checkpointId, dataFile), ++timestamp); + harness.snapshot(checkpointId, ++timestamp); + harness.notifyOfCompletedCheckpoint(checkpointId); + + assertMaxCommittedCheckpointId(jobId, operatorId, checkpointId); + + IcebergFilesCommitter committer = (IcebergFilesCommitter) harness.getOperator(); + + // checkpointId really was committed -- must verify true. + assertThat( + committer.verifyCommitEventuallySucceeded( + jobId.toString(), operatorId.toHexString(), checkpointId, "test")) + .isTrue(); + + // A checkpoint id that was never committed must exhaust the retry budget and report false, + // so commitOperation knows it's not safe to treat this as a no-op and must rethrow. + assertThat( + committer.verifyCommitEventuallySucceeded( + jobId.toString(), operatorId.toHexString(), 999L, "test")) + .isFalse(); + } + } + @TestTemplate public void testCommitTxn() throws Exception { // Test with 3 continues checkpoints: