Skip to content

[spark] Commit a streaming micro-batch through the stream write API - #10105

Open
zhuxiangyi wants to merge 5 commits into
apache:masterfrom
zhuxiangyi:spark-streaming-write-path
Open

zhuxiangyi wants to merge 5 commits into
apache:masterfrom
zhuxiangyi:spark-streaming-write-path

Conversation

@zhuxiangyi

@zhuxiangyi zhuxiangyi commented Sep 22, 2026 •

Copy link
Copy Markdown
Contributor

Purpose

Closes #10010. Fixes #9666. This is the implementation #9667 was asked to be compared against, so #9667 stays
open and untouched while the two are weighed; whichever lands, the other is closed.

A Spark Structured Streaming write goes through the batch write API: every micro-batch builds its
own write builder and commits through a one-shot committer that it closes again. Such a committer
cannot recognise a micro-batch a previous run already committed, so a query that fails between the
sink returning from addBatch and Spark recording that batch in its commit log writes the whole
batch a second time on restart (#9666).

This gives the sink the shape a Flink job has:

  • one commit user per query, stable across its runs, derived from the query id Spark persists
    in the checkpoint — new when a checkpoint is recreated, unchanged when a query resumes from one.
    Batch ids are only unique within one checkpoint, so the identity must not outlive it and is not
    configurable. commit.user-prefix is read from the table as usual and prefixes the derived user,
    as it does the random user of a batch write, so a job upgrading to this keeps the name its table
    was configured with.
  • one StreamTableCommit per run, created with the first micro-batch and closed when the query
    terminates, so the maintenance a commit starts (tag creation, partition and snapshot expiration)
    is not cancelled by a close after every batch.
  • micro-batch n committed under identifier n + 1, the way Flink numbers its checkpoints,
    which is also what a compacted-full scan recognises a scheduled full compaction by. The
    executors write it with prepareCommit(waitCompaction, identifier).
  • the first micro-batch of a run goes through filterAndCommit, the ones after it through
    commit
    — the division of work a Flink committer makes between restoring and its steady state.
  • a write to a postpone bucket table writes the postpone bucket, as FlinkSinkBuilder does for
    a stream, instead of the fixed-bucket paths meant for a batch job that ends with its commit. Those
    rows become readable once a compaction has sorted them into real buckets. This also closes the
    hole where such a write went through the staged committer, which cannot skip a replay at all.
  • a replay that can no longer be recognised is refused. A replay is recognised by the snapshots
    its commit user left behind, which expiration eventually removes. Before each commit the sink
    records under the checkpoint location which snapshot was the latest; if a query comes back after
    expiration has gone past it, the query fails, naming that marker, instead of possibly committing
    the micro-batch twice. The checkpoint is the writer's checkpointLocation option or, when it was
    configured through the session and a query name, the one Spark resolved for the running query; a
    query inside a stream execution whose checkpoint cannot be found fails rather than write without a
    marker.

Relation to #9667

#9667 keeps the one-shot batch committer and makes it drive a filtered commit, which needed three
switches on InnerTableCommit (checkFilesExistence, checkAppendFiles, inlineMaintenance) plus
filterCommittedIgnoresLastSafeSnapshot and a withCommitUser on two write builders. None of that
exists here: a long-lived stream committer needs no switch, because it does not file-list what it
just wrote on a steady-state commit, its maintenance has a next commit to report to, and it never
recomputes a last-safe-snapshot bound.

The only core change left is the chain-table overwrite callback (first commit), which is
independent of either approach: an overwrite of a chain table publishes its delta snapshot and only
then clears the snapshot-branch files it supersedes; when that cleanup fails after the snapshot is
published, the retried commit is recognised as a replay and its retry has to redo the cleanup for
exactly what that overwrite superseded: the files the snapshot branch held then, followed through
the compactions of the snapshot branch since. Where the snapshot branch stood is recorded as the
overwrite is published, in a property of the overwrite's snapshot, through a new
CommitCallback#beforeOverwrite hook whose default does nothing: the two branches commit
independently, so their commit times cannot tell which snapshot branch commits came before the
overwrite. When the cleanup cannot be done exactly — a compaction merged the superseded rows with
later ones, another commit such as a rescale rewrote them, the recorded snapshot has expired, or the
overwrite recorded none — the retry fails instead of reporting a cleanup it did not do.

Tests

PaimonSinkIdempotencyTest (26 cases) and PaimonSinkTest, plus ChainTableFileStoreTableTest and
SimpleTableTestBase in core. Most of the sink cases were written for #9667 against externally
visible behaviour, so they carry over unchanged; what changed:

  • complete mode replay on a postpone bucket table ... now asserts the rows wait in bucket -2 and
    reads them after sys.compact;
  • one committer commits every micro-batch of a query and closes with it (new) counts committers
    through a commit.callbacks implementation: one per query, closed when the query stops;
  • maintenance of a micro-batch runs while the query is alive replaces the case that asserted
    maintenance had to finish before a per-batch committer closed;
  • in core, testFilterAndCommitAcrossRunsOfOneCommitUser replaces the two cases covering the
    removed switches: one committer commits two identifiers, a restarted run replays the last one and
    commits the next.
  • review round (reproduced first, each fails without its fix): a new checkpoint must not skip its batches under a named commit user, a replay whose snapshot may have expired is refused, not committed twice with the replay check refuses only what it cannot tell for the cases it must
    let through, and in core testChainOverwriteReplayAfterSnapshotBranchCompaction,
    testChainOverwriteReplayRefusesWhenCompactionMergedLaterData,
    testChainOverwriteReplayRefusesAfterSnapshotBranchRewrite,
    testChainOverwriteReplayAfterOverwriteTimeSnapshotExpired,
    testChainOverwriteReplayWithSnapshotBranchCommitAtTheSameTime,
    testChainOverwriteReplayRefusesWithoutRecordedPosition,
    testChainOverwriteReplayWithEmptySnapshotBranchKeepsLaterData and
    testChainOverwriteFailsBeforePublishingWhenSnapshotBranchIsUnreadable. The chain tests now fail
    the snapshot branch cleanup by making its snapshot directory read-only, since the overwrite reads
    where that branch stands before publishing.
  • a replay into every kind of table the sink writes to — dynamic bucket, cross partition, postpone
    bucket with real buckets, deletion vector, lookup changelog, row tracking — creates no snapshot;
    a batch that failed before its commit is committed once on restart, and the failed run closes its
    committer; a query resuming a checkpoint written by the previous per-batch commits keeps
    committing, on a table whose older snapshots have expired.
  • the replay refusal also holds for a checkpoint configured through
    spark.sql.streaming.checkpointLocation and a query name, and a query whose checkpoint cannot be
    found fails closed; a marker that was not written completely, or one on a table without any
    snapshot, lets the batch through, and deleting the marker the refusal names commits it.

Run locally: paimon-core in full, paimon-spark-ut in full on Spark 3.5, and the sink suites on
Spark 3.2 / 3.4 / 3.5 / 4.0 / 4.1. Mutation checks: closing the committer after every micro-batch,
never closing it on termination, never filtering a possible replay, never refusing one that cannot
be told, refusing one that can, and each of the three changes to the chain table retry all fail the
cases that cover them.

Postpone write performance (asked for in #9667), local, one JVM, 5 micro-batches of 50k rows
into a primary-key postpone table with postpone.default-bucket-num = 4, two runs:

per micro-batch (steady state) compaction read after
streaming to bucket -2 (this PR) 188-286 ms 2.1-2.3 s 64-65 ms
fixed-bucket write (before) 293-670 ms - 156-213 ms

Committing a micro-batch gets about twice as cheap; the bucketing it no longer does is what the
compaction then does, as for a Flink streaming job.

API and Format

No format change, no new option. CommitCallback gets a beforeOverwrite method with a default that
does nothing, so existing callbacks are unaffected. Behaviour change: a
streaming write to a postpone bucket table now lands in the postpone bucket and becomes readable
after a compaction, as with Flink.

A query that upgrades to this keeps its Spark checkpoint: the sink stores nothing in it. Its commit
user changes from a random one per micro-batch to the stable one, which has no history yet, so the
first micro-batch after the upgrade is committed as usual, neither filtered nor duplicated.

Documentation

docs/docs/spark/structured-streaming.md: an "Exactly-once" section covering the commit user, the
commit identifier, the committer's lifetime, the replay check and what a postpone bucket table does.

@zhuxiangyi
zhuxiangyi force-pushed the spark-streaming-write-path branch from b16121f to 5e96527 Compare September 22, 2026 12:25

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This PR addresses a real end-to-end correctness bug (#9666): a Spark restart can replay a micro-batch that Paimon already committed, duplicating append rows or re-applying aggregation. Reusing the stream write/commit protocol is a sound direction. I ran the Spark 3 PaimonSinkIdempotencyTest suite after packaging the codegen plugin: 17/17 passed. I then ran three additional focused failure-path reproductions. All found production blockers.

[P1] Reusing the documented explicit commit user across a new checkpoint drops new data. The new write.stream.commit-user description and docs/docs/spark/structured-streaming.md:83-86 explicitly recommend pinning the identity when moving to a new checkpoint. A new checkpoint starts Spark batch IDs at 0, so StreamingWrite.commitIdentifier restarts at 1. The first batch of the new query enters filterAndCommit (paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/StreamingWriteContext.scala:70-84), which treats identifier 1 as already committed by the old query and skips its data. I reproduced this by adding the same explicit commit user to both runs of the existing a new query reusing a checkpoint location must not skip its batches test; after resetting the checkpoint, the table contained only (1, old) and was missing (2, new) (16 other cases passed). The temporary test edit was reverted. StreamingWriteContext also advances lastCommitted after a skipped replay, so later batches take the unfiltered fast path with reused identifiers. Please remove the promise that a commit user can be carried to a fresh checkpoint and guard against this destructive configuration, or introduce a distinct epoch/identifier strategy.

[P1] Chain overwrite retry misses old rows after snapshot-branch compaction. paimon-core/src/main/java/org/apache/paimon/metastore/ChainTableOverwriteCommitCallback.java:156-173 intersects file identifiers from the snapshot branch as of the overwrite with its latest file identifiers. If a delta overwrite snapshot is published, its cleanup fails, and the snapshot branch compacts those old files before the retry, all identifiers change. The intersection is empty, so the retry returns without clearing the old rows; reads continue to return those rows instead of the delta overwrite. A focused ChainTableFileStoreTableTest reproduction verified that compaction replaced every old file and then failed the final read assertion: actual old/old-2 rows, expected new. The audit test was reverted. Cleanup needs to track the superseded data across file rewrites, or prevent/resolve this race before treating replay as complete.

[P1] A replay after snapshot expiration is committed again despite the exactly-once claim. FileStoreCommitImpl.java:291-316 deduplicates only against the latest retained snapshot of the commit user; SnapshotManager.java:713-728 stops searching when an older snapshot has expired. I reproduced this in a focused core test: stream user A committed identifier 1 (snapshot 1), user B committed another row (snapshot 2), ExpireSnapshotsImpl.expireUntil(1, 2) removed snapshot 1, and A replayed identifier 1. filterAndCommit returned 1 and published snapshot 3, committing the same batch again; correct replay behavior would return 0 and keep snapshot 2 latest. The audit test was temporary. A production table with aggressive retention or other writers can reach this after the sink has committed but before Spark records the batch in its checkpoint. The documented exactly-once guarantee needs a durable commit-identifier record or an enforceable retention/replay bound; please add a regression test for this interleaving.

CI: Spark 3 and Spark 4 jobs passed. The single failing Flink connector/CDC job failed while starting the Kafka Testcontainers image in KafkaAWSDMSSyncTableActionITCase, outside the changed paths; it should be rerun before merge. GitHub currently reports merge conflicts (mergeStateStatus=DIRTY), which also need resolving and testing. The worktree is clean.

Requirement fit: SUPPORTED. Implementation: BLOCKED on the three reproduced data correctness failures.

@zhuxiangyi
zhuxiangyi force-pushed the spark-streaming-write-path branch from 5e96527 to faa4d5b Compare September 24, 2026 06:46

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Requirement fit: SUPPORTED. #9666 reproduces silent duplicate rows or wrong aggregation results after a Spark micro-batch replay. Reusing the stream protocol has real end-to-end value.

Implementation: FINDINGS. I re-reviewed the current head faa4d5b for production and request changes on four P1 correctness paths:

  1. A fresh checkpoint can silently lose new batches with the documented explicit commit user. paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/StreamingWriteContext.scala:72-84 filters identifier 1 as a replay when an operator reuses write.stream.commit-user for a new checkpoint. The new query starts batch numbering at 0; the old query's identifier 1 makes its first new batch disappear, and subsequent batches can take the unfiltered path. docs/docs/spark/structured-streaming.md:83-86 still recommends pinning the identity across a new checkpoint. My earlier focused reproduction lost the new row; these paths are unchanged on this head. Remove that unsafe recommendation and reject the configuration, or give a fresh checkpoint a distinct identifier epoch.

  2. A replay can commit twice after snapshot expiration. The new filterAndCommit call at StreamingWriteContext.scala:76 depends on retained snapshot history (FileStoreCommitImpl.java:291-316, SnapshotManager.java:713-728). If another writer commits and retention expires this user's committed snapshot before Spark records the batch, a replay finds no prior identifier and publishes it again. My earlier focused reproduction produced snapshot 3 instead of keeping snapshot 2; the relevant code is unchanged. Persist deduplication state independently of expiring snapshots, or define and enforce a retention bound covering the replay window.

  3. Snapshot-branch compaction defeats overwrite cleanup retry. paimon-core/src/main/java/org/apache/paimon/metastore/ChainTableOverwriteCommitCallback.java:157-175 intersects file identifiers from the overwrite-time snapshot with current identifiers. If the delta overwrite publishes, cleanup fails, and the snapshot branch compacts old files before replay, every identifier can change; the retry returns without removing the old rows. An exact-head temporary audit test read old, old-2 instead of delta-new. Track superseded logical data across rewrites, or prevent a retry from silently succeeding when cleanup cannot be proved.

  4. Expiring the overwrite-time snapshot also makes cleanup silently impossible. ChainTableOverwriteCommitCallback.java:152-155 returns when its asOf snapshot is no longer retained, even if the latest snapshot still contains the old rows. An exact-head temporary audit test reproduced old, later after replay where only later should remain. Persist the cleanup intent beyond snapshot retention, or retain the required history and fail explicitly when it is unavailable.

Verification on faa4d5b: ChainTableFileStoreTableTest 49/49 passed after packaging the codegen loader; Spark 3 PaimonSinkIdempotencyTest 17/17 passed. The two additional exact-head audit tests above failed as described, showing gaps in the passing suites. The earlier focused Spark reproductions remain applicable because their affected paths are byte-identical on this head. CI for this revision was still running when reviewed. These data correctness failures block production use.

@zhuxiangyi
zhuxiangyi force-pushed the spark-streaming-write-path branch 2 times, most recently from daaedf9 to b3dcfca Compare September 24, 2026 07:43
@zhuxiangyi

Copy link
Copy Markdown
Contributor Author

Thanks — all four reproduced first as failing tests, then fixed in b3dcfca (rebased on master):

  1. Explicit commit user across a new checkpoint: removed write.stream.commit-user (new in this PR, unreleased). Batch ids are only unique within one checkpoint, so the identity is always the query id and not configurable; commit.user-prefix can still name it. The docs no longer recommend pinning it.
  2. Replay after expiration: before each commit the sink records under the checkpoint location which snapshot was the latest. On a replay with no snapshot of the commit user left, if expiration has gone past that snapshot, the query fails naming the marker instead of committing again; otherwise the lookup is exact (expiration removes the oldest first). Documented as a retention bound.
  3. Compaction before the retry: the retry follows the superseded files through the snapshot branch's compactions since the overwrite and clears their outputs; if a compaction merged them with later rows, it fails instead of reporting success.
  4. Overwrite-time snapshot expired: the retry now fails explicitly when the snapshot branch no longer retains a snapshot from before the overwrite (and only skips when the branch had none then).

Regression tests for each, mutation-checked. I also added a replay test for every kind of table the sink writes to (dynamic bucket, cross partition, postpone with real buckets, DV, lookup changelog, row tracking), a batch failing before its commit, and a query upgraded from the old per-batch commits; sink suites pass on Spark 3.2 / 3.4 / 3.5 / 4.0 / 4.1.

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-review of head b3dcfca022: the previous four findings have concrete fixes and regression coverage. I ran the three new chain-table replay methods locally (3/3 pass), including compaction and expired snapshot-branch history. The Spark 3 suite's existing replay cases also exercised the new checkpoint identity and marker behavior.

[P1] Session-configured checkpoints still duplicate a batch after snapshot expiration. PaimonSink.scala:125-127 creates a CommitMarker only when the writer explicitly passes checkpointLocation as an option. Spark also supports a checkpoint configured through spark.sql.streaming.checkpointLocation plus queryName; this PR already tests that route for ordinary replay, and notes that the location never reaches the sink options. In that route marker=None. Once another writer expires the streaming batch's snapshot, StreamingWriteContext.checkReplayIsRecognisable has neither a retained snapshot of this commit user nor a marker, so it returns and filterAndCommit commits the replay again.

I added a temporary integration test that starts a stream with a session-configured checkpoint, commits (1, stream), inserts (2, batch) with snapshot retention min/max = 1, removes Spark's commit-log entry for batch 0 to simulate a crash after addBatch, and restarts the query. The query completed and the table contained (1, stream) twice plus (2, batch). The same test with an explicit writer checkpointLocation is already covered and refuses the replay. Please obtain the effective checkpoint location or otherwise fail closed when this replay can no longer be recognized; add this session-configuration case to the regression suite.

Verification: the targeted core replay tests passed 3/3. Spark 3 ran 25 cases: 21 passed, the temporary audit case exposed this duplicate, and three postpone-bucket cases were blocked by this isolated clone's stale shaded-Avro artifact (NoSuchMethodError), unrelated to the changed sink logic. The exact-head Java core and interoperability CI are green; Spark 3/4 CI is still pending. This data duplication blocks production use.

@zhuxiangyi
zhuxiangyi force-pushed the spark-streaming-write-path branch from b3dcfca to 0b575c3 Compare September 24, 2026 09:13

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-reviewed head 0b575c3c7f after the session-level checkpoint fix. The previous P1 (a query configured via spark.sql.streaming.checkpointLocation plus queryName had no marker and could commit a replay again after snapshot expiration) is fixed on this head: the sink resolves the running query's checkpoint, writes the marker there, and fails closed if it cannot locate the checkpoint. I also checked the fail-closed and explicit checkpoint paths in the new code. I found no remaining code blocker in the changed path.

Local JDK 17 / Spark 3 verification: PaimonSinkIdempotencyTest 26/26 passed, including the new session-configured checkpoint regression and the previously failing expiration/replay scenario. The initial run was blocked before tests by the local sandbox's socket restriction; the same suite passed when rerun with loopback access. This is an end-to-end streaming test, not only a code inspection.

The GitHub Spark 3/Spark 4 and other affected CI jobs for this new head are still pending at review time. Please use their final results as the merge gate; the earlier changes-requested review should be considered superseded by this re-review once CI passes.

@zhuxiangyi
zhuxiangyi force-pushed the spark-streaming-write-path branch from 0b575c3 to 801e9f8 Compare September 24, 2026 09:43
if (latest == null) {
return;
}
Snapshot asOf = snapshotManager.earlierOrEqualTimeMills(overwrite.timeMillis());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Avoid using the overwrite timestamp as the snapshot-branch boundary.

earlierOrEqualTimeMills can select a snapshot-branch commit made after the delta overwrite. The branches have independent histories, and their wall-clock timestamps have millisecond precision and can differ across hosts. I reproduced this with a focused ChainTableFileStoreTableTest: publish the delta overwrite while its snapshot-branch cleanup fails, append a new row to the snapshot branch, then force only that later snapshot's recorded timeMillis to equal the overwrite's. On replay, this line selects the later snapshot as asOf, and the cleanup deletes the new row along with the superseded row; the read falls back to the delta row. The existing testChainOverwriteReplayClearsOnlyWhatTheOverwriteSuperseded passes with separated timestamps, while this equal-timestamp case fails.

Please use an exact snapshot-branch position associated with the overwrite, or fail closed when that position cannot be established. A wall-clock comparison across branches cannot safely decide which files the overwrite superseded.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 27290b3: the overwrite now records where the snapshot branch stands in a property of its snapshot, through a new CommitCallback#beforeOverwrite hook called before it is published, and retry takes that snapshot as the boundary; the time lookup is gone. It fails closed when no position was recorded or the recorded snapshot has expired, and an unreadable snapshot branch fails the overwrite before publishing. Your equal-timestamp case is testChainOverwriteReplayWithSnapshotBranchCommitAtTheSameTime. Details in #10105 (comment).

…is retried

An overwrite of a chain table publishes its delta snapshot and only then
clears, on the snapshot branch, the files the overwrite superseded. When
that cleanup fails after the snapshot is published, the commit is retried
as a replay: filterCommitted recognises the identifier and hands the
committable to the callback's retry, which did nothing, so the superseded
rows stayed on the snapshot branch for good.

retry now clears what that overwrite superseded in its partitions: the
files the snapshot branch held when the overwrite was published, followed
through the compactions of the snapshot branch since, which rewrite them
into new files. Data written to the snapshot branch after the overwrite is
left alone.

Where the snapshot branch stood is recorded as the overwrite is published,
in a property of the overwrite's snapshot: the two branches commit
independently and their clocks cannot order one against the other, so a
lookup by commit time can take a later snapshot branch commit for the state
the overwrite superseded. CommitCallback gets a beforeOverwrite hook for
this, with a default that does nothing. A snapshot branch that cannot be
read fails the overwrite before it is published.

When the cleanup cannot be done exactly, the retry fails instead of
reporting a cleanup it did not do: when a compaction has merged the
superseded rows with rows written after the overwrite, when another commit
such as a rescale has rewritten them, when the snapshot branch no longer
retains the snapshot the overwrite recorded, or when the overwrite recorded
none.
A Structured Streaming write went through the batch write API: every
micro-batch built its own write builder and committed through a one-shot
committer that it closed again. Such a committer cannot recognise a
micro-batch a previous run already committed, so a query that failed
between the sink returning from addBatch and Spark recording that batch in
its commit log wrote the whole batch a second time on restart.

The sink now has the shape a Flink job has:

  - one commit user per query, stable across its runs. It is derived from
    the query id Spark persists in the checkpoint, which is new when a
    checkpoint is recreated and unchanged when a query resumes from one:
    batch ids are only unique within one checkpoint, so the identity must
    not outlive it, and it is not configurable. 'commit.user-prefix'
    prefixes it, as it does the random user of a batch write.
  - one StreamTableCommit per run, created with the first micro-batch and
    closed when the query terminates, so that the maintenance a commit
    starts is not cancelled by a close after every batch.
  - micro-batch n committed under identifier n + 1, the way Flink numbers
    its checkpoints, which is also what a compacted-full scan recognises a
    scheduled full compaction by. The executors write it with
    prepareCommit(waitCompaction, identifier).
  - the first micro-batch of a run goes through filterAndCommit, the ones
    after it through commit: the same division of work a Flink committer
    makes between restoring and its steady state.
  - a write to a postpone bucket table writes the postpone bucket, as
    FlinkSinkBuilder does for a stream, instead of the fixed-bucket paths
    meant for a batch job that ends with its commit. Those rows become
    readable once a compaction has sorted them into real buckets.

A replay is recognised by the snapshots its commit user left behind, which
snapshot expiration eventually removes. Before each commit the sink records
under the checkpoint location which snapshot was the latest. If a query
comes back after expiration has gone past that snapshot, whether the
micro-batch was committed can no longer be told, and the query fails with
the marker to delete once the table shows it was not, instead of possibly
committing it twice.
@zhuxiangyi
zhuxiangyi force-pushed the spark-streaming-write-path branch from 801e9f8 to 27290b3 Compare September 24, 2026 14:01
@zhuxiangyi

Copy link
Copy Markdown
Contributor Author

Thanks for confirming the session-checkpoint fix. On the timestamp boundary: reproduced first (testChainOverwriteReplayWithSnapshotBranchCommitAtTheSameTime, a later snapshot-branch commit carrying the overwrite's millisecond had its new row deleted), then fixed in 27290b3 by recording the exact position:

  • CommitCallback gets a beforeOverwrite(ManifestCommittable) hook (default no-op), called once before an overwrite is published. ChainTableOverwriteCommitCallback uses it to record the snapshot branch's latest snapshot id (0 if none) in a property of the delta overwrite's snapshot, so it is published atomically with it.
  • retry takes that snapshot as the boundary; no time comparison is left. It fails closed when the overwrite recorded no position (e.g. published before this change) or when the recorded snapshot has expired, and does nothing when the branch was empty.
  • If the snapshot branch cannot be read, the overwrite now fails before it is published, so there is nothing for a replay to redo.
  • One inherent limit, noted in the code: a snapshot-branch commit landing between reading that id and publishing the delta snapshot (a few milliseconds) counts as later and is kept; the existing call path's truncate has the same window.

Also, while auditing coverage I found and fixed one more case of the kind you flagged for compaction: an overwrite-kind rewrite of the snapshot branch (e.g. a rescale) before the retry let it succeed with the superseded rows still there; such rewrites now make the retry fail explicitly (testChainOverwriteReplayRefusesAfterSnapshotBranchRewrite).

Tests: new cases for the missing position, the empty branch and an unreadable branch; the chain tests now fail the cleanup by making the snapshot branch read-only instead of pointing it at a missing branch, since the overwrite reads that branch first. All mutation-checked. Core is green apart from the Docker-based PostgresqlCatalogTest; SparkChainTableITCase, CompactChainTableProcedureTest, the Flink chain ITCases and the sink suites on Spark 3.2 / 3.4 / 3.5 / 4.0 / 4.1 pass. The two failures on the previous CI run were a 403 from Maven Central and a MySQL CDC timeout, unrelated to this change.

@JingsongLi

Copy link
Copy Markdown
Contributor

Re-reviewed latest head 27290b3 for production, focusing on the new chain-table overwrite replay path after the earlier Spark streaming idempotency reviews. The previous equal-timestamp cross-branch bug is addressed: the delta overwrite now records the snapshot branch's exact snapshot ID before publication, and retry uses that position instead of comparing independent branch clocks. The retry follows subsequent file rewrites and fails closed when the position has expired or old/new data can no longer be separated.

Local JDK 8 ChainTableFileStoreTableTest passed 57/57, including failed cleanup, later snapshot-branch writes, compaction, the equal-timestamp regression, and missing-position cases. The Spark 3 and Spark 4 CI jobs pass on this head. Core CI is red only because S3FileIOTest could not pull its MinIO Docker image (both JDK 8 and 11), not because a changed-path test failed; please rerun those jobs before merge. git diff --check passed. I found no new code blocker in this follow-up.

@zhuxiangyi

Copy link
Copy Markdown
Contributor Author

Fixed in 1e664d7d4:

  • Preserve the original commit marker across retries and fail closed on corrupt markers.
  • Validate replay safety after filtering to avoid races with snapshot expiration.

Added regression coverage for repeated recovery, interrupted marker writes, corrupt markers, and concurrent expiration (including stale hints). All 122 targeted tests pass.

@JingsongLi

Copy link
Copy Markdown
Contributor

I reviewed the stream write/commit path, checkpoint recovery, and chain-table overwrite cleanup on the current head 1e664d7d4a03dc4f5510f45d62bbb32c5f99c6a9. The change has clear end-to-end value, but two recovery/lifetime cases need addressing before relying on its exactly-once behavior in production.

  1. [P2] Freeze the commit identity across changes to commit.user-prefix. PaimonSink.withPrefix recomputes the complete commit user from the current table/write options on each run. Using a real Spark streaming query, I committed one append row, stopped the query, removed checkpoint/commits/0 to reproduce the acknowledged-commit crash window, changed the table property from old-job to new-job, and resumed the same checkpoint. Spark preserved the query UUID, but the commit users became old-job_spark-query-... and new-job_spark-query-...; the table contained two identical rows instead of one. The marker is also treated as belonging to another user and replaced, so it does not fail closed. Persist and reuse the full identity for the checkpoint, or explicitly reject a prefix change when resuming it. Please cover both a table-property change and a write-option change.

  2. [P2] Match the terminating query run, not just the persistent query ID. registerCloseOnTermination checks only event.id, although a restart keeps that ID and changes runId. Query termination is dispatched asynchronously. I blocked a real Spark listener's progress callback, stopped the first run, resumed its checkpoint and committed the next batch, then released the queued termination event. Both committers were closed while the restarted query was still active (created=2, closed=2, live=true), and the new run's listener removed itself. Another micro-batch recreated a third committer. After stopping the restarted query and observing its actual termination, that committer remained unclosed (created=3, closed=2, live=false); all three committed rows were readable. Bind the listener to the current run ID and ensure it only closes that run's already-created context; add a regression with delayed termination from the preceding run.

Validation: 82 selected Core tests and all 41 sink tests on each of Spark 3.5.8 and Spark 4.1.2 passed with normal Maven checks. An independent real chain-table probe confirmed that initial cleanup and replay preserve a snapshot baseline needed by a following delta partition. The two additional actual-streaming probes above fail on this head. The existing Flink 1 Common CI job was cancelled after approximately 100 minutes, so that part of CI still needs a completed run.

…ed run

A checkpoint resumed under a different 'commit.user-prefix' derived a new
commit user, so the replay of a batch the previous run committed was not
recognised and was committed twice. The marker the previous run left names
its commit user; a run with the same identity but another prefix is now
refused, for a prefix set as a table property or a write option.

A restarted query keeps its id and Spark delivers terminations
asynchronously, so a late termination of the previous run closed the
committer of the restarted one and left the committer it recreated
unclosed. The termination listener now matches the run id, and closes the
context only if the run created one.
@zhuxiangyi

Copy link
Copy Markdown
Contributor Author

Thanks for the review. Both issues are reproduced and fixed in e40250f:

  • P2-1: Resuming a checkpoint with a different commit.user-prefix now fails with a clear error. Tests cover a table-property change and a write-option change.
  • P2-2: The termination listener now matches the run ID and only closes the context that run created. A regression test covers a delayed termination from the previous run.

The sink tests pass on Spark 3 and Spark 4. I'll make sure the Flink CI job completes.

@JingsongLi

Copy link
Copy Markdown
Contributor

@zhuxiangyi Given the complexity of these changes, just how important is this PR? Are you contributing it for use in a production application? What specific use case is it intended for?

@zhuxiangyi

Copy link
Copy Markdown
Contributor Author

Thanks for asking. We plan to move an ingestion pipeline for log tables (append-only, no primary key) from Flink CDC → Paimon to a long-running Spark Structured Streaming job that reads JSON from Kafka and writes to Paimon through writeStream.format("paimon"). We found this while evaluating the migration: if the driver dies after a micro-batch is committed but before Spark records it in its checkpoint (an OOM, a lost node, the cluster manager restarting it, or a redeploy that kills the job), the restarted query replays that batch and Paimon commits it again (#9666). With no primary key to merge them, every row of that batch is then duplicated in the log table. On recovery, Paimon's Flink committer filters what a restored checkpoint already committed; the Spark sink commits every micro-batch under a random user, so it has nothing to filter by.

So what we need is exactly-once for the Spark sink, on par with Flink. I understand the concern about the size of the change. If it helps review, I can split it: the chain-table overwrite cleanup on retry is independent of the Spark sink and can be its own PR, leaving this one with the stream write/commit path and the replay check.

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

Labels

None yet

Projects

None yet

2 participants