-
Notifications
You must be signed in to change notification settings - Fork 1
Flink 1.18: harden IcebergFilesCommitter against ambiguous-commit duplicates #5
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: 1.8.1-affirm-patch
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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; | ||
|
|
@@ -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,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"; | ||
| 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; | ||
|
|
@@ -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 | ||
|
|
@@ -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(); | ||
|
|
@@ -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); | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 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())) { | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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( | ||
| "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) { | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) {
Suggest filtering the walk to snapshots that are actually Flink commits (i.e. only increment |
||
| 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, | ||
|
|
@@ -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", | ||
|
|
@@ -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); | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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);
|
||
| } 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; | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 |
||
| } | ||
|
|
||
| return false; | ||
| } | ||
|
|
||
| @Override | ||
| public void processElement(StreamRecord<FlinkWriteResult> element) { | ||
| FlinkWriteResult flinkWriteResult = element.getValue(); | ||
|
|
||
There was a problem hiding this comment.
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
No test sets this property to a non-default value or otherwise probes the truncation boundary in
collectRecentlyCommittedFilePaths'sinspected < recentSnapshotLookbackloop (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.