Skip to content

[Data] Support commit-aware checkpointing for transactional Datasinks (Paimon) #66491

Description

@yujiabin123

Description

Ray Data checkpoints non-_FileDatasink sinks right after Datasink.write() returns, before the sink's own transaction commits. For Paimon this is unsafe:

  • each write task returns CommitMessages and its row checkpoint is persisted immediately (post-write transform, same task);
  • the table snapshot is committed only once, in on_write_complete() on the driver, after ALL tasks finish.

So any failure after the first task completes — a later task crash, worker loss, or the final commit itself — leaves rows checkpointed but uncommitted. On resume, filter_rows_for_block skips those rows at the source: they are never rewritten and never appear in any snapshot. Enabling checkpointing here is worse than not enabling it — without it, a full rerun rewrites everything into one final commit and converges correctly.

(For _FileDatasink this is handled correctly via pending/committed 2PC; non-file sinks get a runtime warning acknowledging at-least-once semantics.)

Use case

Long-running backfills/upserts into Paimon primary-key tables need resumable writes. Desired invariant: commit first, then publish the input checkpoint — persist retryable write results, commit them at the sink's transaction boundary, and publish row checkpoints only after that succeeds.

Iceberg is getting this via the #66387–#66394 series (persisted write-result metadata; row checkpoint published last as the completion marker; recovery keyed on a generation marker in snapshot properties). Two Paimon-specific questions for that pattern:

  • CommitMessages are in-memory objects — is a serializable write-result artifact (as Iceberg does for DataFile metadata) the expected prerequisite?
  • Iceberg's current support is APPEND-only; Paimon primary-key upserts are the main use case here — should upsert semantics be in scope, or follow-up work?

Happy to contribute a PR following the Iceberg checkpoint-datasink wrapper structure, with maintainer guidance on the above.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions