zhuxiangyi opened a new pull request, #10105:
URL: https://github.com/apache/paimon/pull/10105

   ### Purpose
   
   Closes #10010. 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 —
     or set explicitly with the new `write.stream.commit-user` option. It is 
read only from sources
     scoped to one query (the writer's options, the session conf), never from a 
table property that
     every writer of the table would share.
   - **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.
   
   ### 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 the files that overwrite superseded.
   
   ### Tests
   
   `PaimonSinkIdempotencyTest` (15 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.
   
   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, and never filtering a possible replay each 
fail the cases that
   cover them (3, 1 and 9 cases).
   
   **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. New Spark connector option `write.stream.commit-user`. 
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.
   
   ### Documentation
   
   `docs/docs/spark/structured-streaming.md`: an "Exactly-once" section 
covering the commit user, the
   commit identifier, the committer's lifetime and what a postpone bucket table 
does.
   


-- 
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]

Reply via email to