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]

Reply via email to