Repository navigation
[spark] Commit a streaming micro-batch through the stream write API - #10105
zhuxiangyi wants to merge 5 commits into
Conversation
b16121f to
5e96527
Compare
JingsongLi
left a comment
There was a problem hiding this comment.
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.
5e96527 to
faa4d5b
Compare
JingsongLi
left a comment
There was a problem hiding this comment.
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:
-
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-84filters identifier 1 as a replay when an operator reuseswrite.stream.commit-userfor 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-86still 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. -
A replay can commit twice after snapshot expiration. The new
filterAndCommitcall atStreamingWriteContext.scala:76depends 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. -
Snapshot-branch compaction defeats overwrite cleanup retry.
paimon-core/src/main/java/org/apache/paimon/metastore/ChainTableOverwriteCommitCallback.java:157-175intersects 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 readold, old-2instead ofdelta-new. Track superseded logical data across rewrites, or prevent a retry from silently succeeding when cleanup cannot be proved. -
Expiring the overwrite-time snapshot also makes cleanup silently impossible.
ChainTableOverwriteCommitCallback.java:152-155returns when itsasOfsnapshot is no longer retained, even if the latest snapshot still contains the old rows. An exact-head temporary audit test reproducedold, laterafter replay where onlylatershould 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.
daaedf9 to
b3dcfca
Compare
|
Thanks — all four reproduced first as failing tests, then fixed in b3dcfca (rebased on master):
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
left a comment
There was a problem hiding this comment.
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.
b3dcfca to
0b575c3
Compare
JingsongLi
left a comment
There was a problem hiding this comment.
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.
0b575c3 to
801e9f8
Compare
| if (latest == null) { | ||
| return; | ||
| } | ||
| Snapshot asOf = snapshotManager.earlierOrEqualTimeMills(overwrite.timeMillis()); |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
801e9f8 to
27290b3
Compare
|
Thanks for confirming the session-checkpoint fix. On the timestamp boundary: reproduced first (
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 ( 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 |
|
Re-reviewed latest head Local JDK 8 |
|
Fixed in
Added regression coverage for repeated recovery, interrupted marker writes, corrupt markers, and concurrent expiration (including stale hints). All 122 targeted tests pass. |
|
I reviewed the stream write/commit path, checkpoint recovery, and chain-table overwrite cleanup on the current head
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.
|
Thanks for the review. Both issues are reproduced and fixed in e40250f:
The sink tests pass on Spark 3 and Spark 4. I'll make sure the Flink CI job completes. |
|
@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? |
|
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 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. |
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
addBatchand Spark recording that batch in its commit log writes the wholebatch a second time on restart (#9666).
This gives the sink the shape a Flink job has:
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-prefixis 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.
StreamTableCommitper run, created with the first micro-batch and closed when the queryterminates, so the maintenance a commit starts (tag creation, partition and snapshot expiration)
is not cancelled by a close after every batch.
ncommitted under identifiern + 1, the way Flink numbers its checkpoints,which is also what a
compacted-fullscan recognises a scheduled full compaction by. Theexecutors write it with
prepareCommit(waitCompaction, identifier).filterAndCommit, the ones after it throughcommit— the division of work a Flink committer makes between restoring and its steady state.FlinkSinkBuilderdoes fora 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.
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
checkpointLocationoption or, when it wasconfigured 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) plusfilterCommittedIgnoresLastSafeSnapshotand awithCommitUseron two write builders. None of thatexists 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
retryhas to redo the cleanup forexactly 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#beforeOverwritehook whose default does nothing: the two branches commitindependently, 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) andPaimonSinkTest, plusChainTableFileStoreTableTestandSimpleTableTestBasein core. Most of the sink cases were written for #9667 against externallyvisible behaviour, so they carry over unchanged; what changed:
complete mode replay on a postpone bucket table ...now asserts the rows wait in bucket -2 andreads them after
sys.compact;one committer commits every micro-batch of a query and closes with it(new) counts committersthrough a
commit.callbacksimplementation: one per query, closed when the query stops;maintenance of a micro-batch runs while the query is alivereplaces the case that assertedmaintenance had to finish before a per-batch committer closed;
testFilterAndCommitAcrossRunsOfOneCommitUserreplaces the two cases covering theremoved switches: one committer commits two identifiers, a restarted run replays the last one and
commits the next.
a new checkpoint must not skip its batches under a named commit user,a replay whose snapshot may have expired is refused, not committed twicewiththe replay check refuses only what it cannot tellfor the cases it mustlet through, and in core
testChainOverwriteReplayAfterSnapshotBranchCompaction,testChainOverwriteReplayRefusesWhenCompactionMergedLaterData,testChainOverwriteReplayRefusesAfterSnapshotBranchRewrite,testChainOverwriteReplayAfterOverwriteTimeSnapshotExpired,testChainOverwriteReplayWithSnapshotBranchCommitAtTheSameTime,testChainOverwriteReplayRefusesWithoutRecordedPosition,testChainOverwriteReplayWithEmptySnapshotBranchKeepsLaterDataandtestChainOverwriteFailsBeforePublishingWhenSnapshotBranchIsUnreadable. The chain tests now failthe snapshot branch cleanup by making its snapshot directory read-only, since the overwrite reads
where that branch stands before publishing.
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.
spark.sql.streaming.checkpointLocationand a query name, and a query whose checkpoint cannot befound 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-corein full,paimon-spark-utin full on Spark 3.5, and the sink suites onSpark 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: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.
CommitCallbackgets abeforeOverwritemethod with a default thatdoes 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, thecommit identifier, the committer's lifetime, the replay check and what a postpone bucket table does.