rayhondo opened a new issue, #18402: URL: https://github.com/apache/iceberg/issues/18402
### Apache Iceberg version 1.11.0, plus upstream cherry-picks of #16648 (the `flink/v2.0` port of #16011), #17194, #18101 and #17437. The validator from #14517/#14637 is already in stock 1.11.0. None of the cherry-picks change the validator, `SnapshotProducer`, `MergingSnapshotProducer` or the REST client. Every line cited below is from the stock `apache-iceberg-1.11.0` tag. ### Query engine Flink (`DynamicIcebergSink` / `DynamicCommitter`, `flink/v2.0`), REST catalog ### Please describe the bug 🐞 A REST catalog applied a `RowDelta` commit from `DynamicCommitter`, but the client got a 409 back for the same commit. On the retry, `MaxCommittedCheckpointIdValidator` (added in #14517) found the checkpoint already on `main` and threw `MaxCommittedCheckpointMismatchException`. That exception extends `ValidationException`, which implements `CleanableFailure`, so `SnapshotProducer.commit()` ran `cleanAll()`. `cleanAll()` deleted the manifest list and the new data manifest of the snapshot the catalog had already made current. `DynamicCommitter` then caught the exception and logged `Skipping commit operation rowDelta ... already contains changes for checkpoint N`. The table was left with a current snapshot whose manifest list returned 404. Every later commit and committer restore failed with `NoSuchKey` until the two objects were restored from S3 object versions. Observed sequence (single writer, no other snapshot between S(n-1) and S(n+1), no compaction in the window): 1. The committer builds snapshot S(n) for checkpoint N with parent S(n-1) and POSTs it with requirement `assert-ref-snapshot-id main == S(n-1)`. 2. The catalog applies it, and `main` moves to S(n). The manifest list `snap-S(n)-1-<uuid>.avro` and manifest `<uuid>-m0.avro` exist in storage. 3. About 64s later the client logs the first and only retry: `Tasks - Retrying task after failure: sleepTimeMs=108 Commit failed: Requirement failed: branch main has changed: expected id S(n-1) != S(n)`. The "actual" id in that message is the snapshot this same commit created. 4. Attempt 2: `apply()` refreshes the table, and the validator sees checkpoint N on `main` and throws `MaxCommittedCheckpointMismatchException: Table already contains staged changes.` 5. About 14s later the manifest list and manifest of S(n) are deleted. Both carry the producer's commit UUID. 6. The committer restarts and fails on every restore with S3 404 until the objects are restored. The client logged no `CommitStateUnknownException`. ### Expected behavior If the table already contains this producer's own snapshot, the retry should end as a success, or at least skip deleting files. Failure cleanup should never delete files referenced by a snapshot that is on the table. ### Root cause (apache-iceberg-1.11.0, `6976e020`) - `core/src/main/java/org/apache/iceberg/UpdateRequirement.java:117-120`: the server builds `... has changed: expected id %s != %s` and throws `CommitFailedException`. - `core/src/main/java/org/apache/iceberg/rest/ErrorHandlers.java:127-128`: a 409 becomes `CommitFailedException`. 500/502/503/504 become `CommitStateUnknownException` (`:129-134`). - `core/src/main/java/org/apache/iceberg/SnapshotProducer.java:471`: `Tasks ... onlyRetryOn(CommitFailedException.class)`. Each attempt calls `apply()` (`:475`), which calls `refresh()` (`:273`) and `runValidations()` (`:279`, `:359-361`). - `SnapshotProducer.java:665-669`: `snapshotId()` is assigned once, so attempt 2 re-stages the same id that is already on `main`. - `flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/DynamicCommitter.java:308`: `MaxCommittedCheckpointMismatchException extends ValidationException`. The validator throws it at `:330-331`. It is registered at `:363-364`, and the catch that logs "Skipping commit" (`:369-379`) wraps `operation.commit()` (`:368`), so it runs after the cleanup. - `api/src/main/java/org/apache/iceberg/exceptions/ValidationException.java:35`: `implements CleanableFailure`. - `SnapshotProducer.java:504-509`: `CommitStateUnknownException` is rethrown without cleanup. Any other `RuntimeException` that is a `CleanableFailure` goes to `Exceptions.suppressAndThrow(e, this::cleanAll)`. - `SnapshotProducer.java:578-583`: `cleanAll()` deletes every entry in `manifestLists`, including attempt 1's list, which S(n) references. `cleanUncommitted(EMPTY_SET)` then reaches `MergingSnapshotProducer.java:1076-1078` and deletes the cached new data and delete manifests. - The success path (`SnapshotProducer.java:521-532`) already refreshes, loads the committed snapshot by id and spares its manifests and manifest list. The failure path has no equivalent check. ### Why 1.10.2 does not hit this (derived from source, not tested) 1.10.2 has no `MaxCommittedCheckpointIdValidator`. `DynamicCommitter.commitOperation` calls `operation.commit()` directly (`DynamicCommitter.java:359` at `apache-iceberg-1.10.2`). On the same retry, `apply()` succeeds. Because `base.snapshot(snapshotId) != null`, the producer takes the rollback branch (`SnapshotProducer.java:443-445`), `updated.changes().isEmpty()` returns (`:453-458`), and the success-path cleanup keeps S(n)'s files. The cleanup code is the same in both versions. 1.11.0 adds a `CleanableFailure` that fires exactly when the commit has already been applied. ### Candidate fix directions (no recommendation yet) 1. **Guard the failure cleanup in `SnapshotProducer.commit()`.** Before `cleanAll()`, refresh, and if `newSnapshotId` is on the table, clean up as the success path does. This covers any `CleanableFailure` after an applied attempt. The cost is an extra `refresh()` on the failure path, and a failure to refresh still needs a safe default. 2. **Make the validator failure non-cleanable**, or have `DynamicCommitter` validate on a refreshed table before `commit()`. This is narrow and local to Flink, but leaves the core gap open for other validators. It also changes `ValidationException` semantics if done in core. 3. **Treat a 409 whose current ref equals the producer's own snapshot as success**, similar to the `CommitStateUnknownException` reconcile in `RESTTableOperations.java:207-243`. This is REST-only and fixes the cause rather than the cleanup, but it does not help other catalogs. Related: #14425 / #14517 (the validator), #17901 (open, adds concurrent-commit validation to `IcebergSink`; it may meet the same cleanup path, unverified). ### Willingness to contribute - [x] I can contribute a fix for this bug independently ### Unverified - What evaluated the commit a second time. Candidates: a catalog-internal retry that re-checked the requirement after persisting, or a resend of the same POST (the client retries POST on 429 or on 503 with `Retry-After`, `ExponentialHttpRequestRetryStrategy.java:91,142-155`; or a proxy retry). We have no catalog-side request log yet. - Why the response took about 64s. - Whether the catalog advertised `idempotency-key-lifetime`. The client sends `Idempotency-Key` only if it does (`RESTSessionCatalog.java:220-222`). - The 1.10.2 behaviour above comes from reading the source; we did not run it. - That the deletes came from `cleanAll()` is inferred from log order and object timestamps, not from storage access logs. -- 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]
