[Bug] IcebergCommitCallback crashes the Flink job on CommitStateUnknownException even though the committer already self-heals on the next checkpoint

Ouverte
#8,875 1 commentaire 0 réactions 0 personnes assignées Voir sur GitHub

Personne n'a encore pris cette issue.

Évaluation

Difficulté
4/5
Temps estimé
3-5 jours
Accessibilité débutants
48/100
Type d'issue
Bug
Clarté
Plutôt claire
Activité
Calme
Stack technique
java

Piste de recherche

Commencez par lire IcebergCommitCallback.call/createMetadata, IcebergRestMetadataCommitter.commitMetadataImpl et FileStoreCommitImpl.tryCommitOnce, puis reproduisez le timeout du REST-catalog avec Flink 2.0. Suivez la manière dont CommitterOperator.notifyCheckpointComplete gère l'exception et ajoutez une couverture pour le comportement prévu au prochain checkpoint, la finalisation étant définie par l'accord des maintainers, ainsi qu'un test de régression pour le cas de commit ambigu.

Rédigé par le modèle d'indexation à partir du texte de l'issue.

Description

Search before asking
  • I searched in the issues and found nothing similar.
Paimon version

1.4.1 (also checked: the same propagation/commit-callback behavior is present on current master)

Compute Engine

Flink 2.0, Paimon table with Iceberg metadata sync enabled against a REST catalog (metadata.iceberg.storage = rest-catalog)

Minimal reproduce step
  1. Configure a Paimon table with Iceberg REST metadata sync enabled, written to continuously by a Flink streaming job (one Iceberg commit attempt per checkpoint).

  2. Cause the Iceberg REST catalog's commit endpoint to respond with a 500/timeout on a single commit POST. The Iceberg REST client surfaces this as org.apache.iceberg.exceptions.CommitStateUnknownException.

  3. This exception propagates unhandled out of IcebergRestMetadataCommitter.commitMetadataImplIcebergCommitCallback.createMetadataWithBase/createMetadataIcebergCommitCallback.callFileStoreCommitImpl.tryCommitOnce → ... → CommitterOperator.notifyCheckpointComplete. Because this happens inside notifyCheckpointComplete, Flink treats it as fatal and restarts the whole job:

    java.lang.RuntimeException: java.lang.RuntimeException: Fail to commit metadata to rest catalog.
    	at org.apache.paimon.iceberg.IcebergRestMetadataCommitter.commitMetadata(IcebergRestMetadataCommitter.java:127)
    	at org.apache.paimon.iceberg.IcebergCommitCallback.createMetadataWithBase(IcebergCommitCallback.java:667)
    	at org.apache.paimon.iceberg.IcebergCommitCallback.createMetadata(IcebergCommitCallback.java:276)
    	at org.apache.paimon.iceberg.IcebergCommitCallback.call(IcebergCommitCallback.java:224)
    	at org.apache.paimon.operation.FileStoreCommitImpl.lambda$tryCommitOnce$11(FileStoreCommitImpl.java:1048)
    	...
    	at org.apache.paimon.flink.sink.CommitterOperator.commitUpToCheckpoint(CommitterOperator.java:215)
    	at org.apache.paimon.flink.sink.CommitterOperator.notifyCheckpointComplete(CommitterOperator.java:192)
    	...
    Caused by: org.apache.iceberg.exceptions.CommitStateUnknownException: Service failed: 500: timeout exceeded
    Cannot determine whether the commit was successful or not, the underlying data files may or may not be needed. ...
    

    Crucially, this callback runs after Paimon's own core table commit (FileStoreCommitImpl.commit) has already succeeded — the Iceberg REST sync is a best-effort side effect, not part of the table's core commit path or durability guarantee.

  4. On the very next checkpoint's commit attempt, commitMetadataImpl reloads the table fresh from the REST catalog (icebergTable = getTable()restCatalog.loadTable(...), not a cached/stale reference) and calls checkBase():

    return currentMetadata.currentSnapshot().snapshotId()
            == newMetadata.currentSnapshot().snapshotId() - 1;
    

    What happens here depends on whether the ambiguous commit from step 2 actually landed server-side:

    • If it landed, the REST-side current snapshot is now exactly one behind the new local metadata's current snapshot, checkBase() returns true, and the commit proceeds normally via updatesForCorrectBase(..., false) — no recreate needed at all.
    • If it didn't land (or the REST side otherwise ends up more than one snapshot behind — this is also unconditionally true whenever the local baseIcebergMetadata passed in is null, a separate trigger for the same branch), checkBase() returns false, and updatesForIncorrectBase()recreateTable() runs: it drops and recreates the Iceberg table with a single snapshot that correctly represents the current live file set (manifests are built cumulatively, so there's no query-visible data loss).

    In neither branch does whether the job crashed in step 3 change the outcome, since the check is always against live server state.

What doesn't meet your expectations?

Given step 4, it's not obvious to us that crashing the whole Flink job in step 3 helps: if the ambiguous commit actually landed, the next checkpoint would have committed normally with no recreate at all; if it didn't land, the committer self-heals via updatesForIncorrectBase()/recreateTable() on the next attempt regardless of whether the job restarted. In both cases the crash adds a full job restart (TaskManager churn, checkpoint restore, backpressure/lag buildup across every table this job writes, not only the one whose commit was ambiguous) on top of an outcome that seems to resolve itself either way.

We'd be curious whether there's a reason CommitStateUnknownException (and similarly ambiguous/retryable REST-catalog errors) is treated as fatal here rather than caught and logged, letting the next checkpoint's commit attempt resolve things via the existing self-healing path — or whether we're missing a case where crashing is actually necessary. Interested in maintainer/community input on this before we look at a PR.

Are you willing to submit a PR?
  • I'm willing to submit a PR!
Langage dominant
Java
Étoiles
3.4k
Forks
1.4k
Merge moyen
1 j 14 h
PR mergées (30 j)
468

Guide de contribution

Aucun guide de contribution indexé pour ce dépôt

Par où commencer

  1. Lisez l'issue en entier, puis le guide de contribution du projet.
  2. Signalez en commentaire que vous la prenez — cela évite que deux personnes fassent le même travail.
  3. Forkez le dépôt et travaillez sur une branche.
  4. Ouvrez une pull request qui référence le numéro de l'issue.

Autres issues de apache/paimon

Toutes les issues de apache/paimon

Issues similaires

Plus d'issues Java

Recevez les nouvelles issues par e-mail

Un résumé court des issues GitHub adaptées aux débutants.