Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -40,13 +40,16 @@
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;
import org.apache.iceberg.RowDelta;
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;
Expand All @@ -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;
Expand All @@ -80,6 +84,40 @@ class IcebergFilesCommitter extends AbstractStreamOperator<Void>
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.
static final String RECENT_SNAPSHOT_LOOKBACK_PROP = "flink.recent-snapshot-lookback";

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Test coverage: no test for the lookback truncation boundary

static final String RECENT_SNAPSHOT_LOOKBACK_PROP = "flink.recent-snapshot-lookback";

No test sets this property to a non-default value or otherwise probes the truncation boundary in collectRecentlyCommittedFilePaths's inspected < recentSnapshotLookback loop (see the comment on line 421 re: the counter-scoping bug). Worth a regression test once that logic is fixed, so the boundary behavior stays covered going forward.

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
// 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. 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;
private final boolean replacePartitions;
Expand Down Expand Up @@ -108,6 +146,9 @@ class IcebergFilesCommitter extends AbstractStreamOperator<Void>
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
Expand Down Expand Up @@ -155,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();
Expand Down Expand Up @@ -275,13 +329,115 @@ private void commitUpToCheckpoint(
manifests.addAll(deltaManifests.manifests());
}

CommitSummary summary = new CommitSummary(pendingResults);
commitPendingResult(pendingResults, summary, newFlinkJobId, operatorId, checkpointId);
NavigableMap<Long, WriteResult> dedupedResults =
dropAlreadyCommittedFiles(pendingResults, newFlinkJobId, operatorId);

CommitSummary summary = new CommitSummary(dedupedResults);
commitPendingResult(dedupedResults, summary, newFlinkJobId, operatorId, checkpointId);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Correctness: no guard for a checkpoint entry fully emptied by dedup

CommitSummary summary = new CommitSummary(dedupedResults);
commitPendingResult(dedupedResults, summary, newFlinkJobId, operatorId, checkpointId);

If dropAlreadyCommittedFiles empties out both dataFiles() and deleteFiles() for a given checkpoint entry (a full duplicate), dedupedResults still contains that now-empty entry, and it still flows into commitPendingResult (which dispatches to commitDeltaTxn/commitAppendTxn per entry) with no check for emptiness. If the resulting all-empty RowDelta/AppendFiles commit is ever rejected by something other than CommitStateUnknownException, that exception propagates out of commitUpToCheckpoint before pendingMap.clear() runs — leaving a later, genuinely-new entry in the same batch stuck pending for a retry where a since-shifted lookback window could misclassify it.

Suggest skipping the commit entirely (and just clearing the entry) when a checkpoint's dedup result is fully empty.

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 {@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.
*/
@VisibleForTesting
@SuppressWarnings("CollectionUndefinedEquality") // CharSequenceSet defines path equality itself
NavigableMap<Long, WriteResult> dropAlreadyCommittedFiles(
NavigableMap<Long, WriteResult> pendingResults, String newFlinkJobId, String operatorId) {
CharSequenceSet recentlyCommittedPaths =
collectRecentlyCommittedFilePaths(newFlinkJobId, operatorId);
if (recentlyCommittedPaths.isEmpty()) {
return pendingResults;
}

NavigableMap<Long, WriteResult> deduped = Maps.newTreeMap();
for (Map.Entry<Long, WriteResult> 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())) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Maintainability: duplicated data-file / delete-file loops

if (recentlyCommittedPaths.contains(file.path())) {

This loop and its twin for delete files at line 386 are near-identical, with copy-pasted LOG.warn blocks differing only in "data file"/"delete file" wording. Consider extracting a shared helper so a future fix or format change to one doesn't get forgotten in the other.

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 {@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) {
table.refresh();

CharSequenceSet paths = CharSequenceSet.empty();
Snapshot snapshot = table.snapshot(branch);
int inspected = 0;
while (snapshot != null && inspected < recentSnapshotLookback) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Correctness: lookback budget is consumed by unrelated snapshots, not just matches

int inspected = 0;
while (snapshot != null && inspected < recentSnapshotLookback) {

inspected increments on every ancestor snapshot visited (line 435), regardless of whether it matches flinkJobId/operatorId. Unrelated snapshots interleaved on the branch — compaction, dangling-delete removal, any other maintenance operation — burn the default 5-snapshot lookback budget before the walk ever reaches the snapshot that actually holds the ambiguous commit's file paths. If 5+ such snapshots land between the ambiguous commit and its retry, dropAlreadyCommittedFiles won't see the duplicate and will re-append it — silently defeating this PR's own stated purpose for the apache#10765 bug it targets.

Suggest filtering the walk to snapshots that are actually Flink commits (i.e. only increment inspected for snapshots carrying a flink.job-id, or better, only for ones matching this exact flinkJobId/operatorId) so non-Flink maintenance snapshots don't silently consume the budget.

Map<String, String> 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<Long, WriteResult> pendingResults,
CommitSummary summary,
Expand Down Expand Up @@ -412,7 +568,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_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 "
+ "(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,
commitStateUnknownMaxVerifyAttempts,
e);
throw e;
}
long durationMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNano);
LOG.info(
"Committed {} to table: {}, branch: {}, checkpointId {} in {} ms",
Expand All @@ -424,6 +615,53 @@ private void commitOperation(
committerMetrics.commitDuration(durationMs);
}

/**
* 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.
*/
@VisibleForTesting
boolean verifyCommitEventuallySucceeded(
String flinkJobId, String operatorId, long checkpointId, String description) {
long delayMs = commitStateUnknownVerifyInitialDelayMs;
for (int attempt = 1; attempt <= commitStateUnknownMaxVerifyAttempts; attempt++) {
try {
Thread.sleep(delayMs);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reliability: blocks the operator's mailbox thread for up to ~31s per ambiguous commit

Thread.sleep(delayMs);

verifyCommitEventuallySucceeded synchronously blocks the calling thread inside notifyCheckpointComplete for up to the full retry budget (default 5 attempts, 1s→16s backoff ≈ 31s total) per ambiguous commit. Worth checking this against Chrono's actual checkpoint/heartbeat timeout configuration — a long enough block here risks tripping Flink's own timeouts and causing a harsher failure than the one this PR is trying to avoid.

} 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,
commitStateUnknownMaxVerifyAttempts,
observedCheckpointId,
flinkJobId,
operatorId);
if (observedCheckpointId >= checkpointId) {
return true;
}

delayMs *= 2;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Maintainability: reimplements existing backoff utility

delayMs *= 2;

This hand-rolled exponential-backoff loop duplicates org.apache.iceberg.util.Tasks.Builder#exponentialBackoff, already available in this repo (core/src/main/java/org/apache/iceberg/util/Tasks.java:183). Not a correctness issue, but it diverges from an existing, already-tested retry utility rather than reusing it.

}

return false;
}

@Override
public void processElement(StreamRecord<FlinkWriteResult> element) {
FlinkWriteResult flinkWriteResult = element.getValue();
Expand Down
Loading