zhuxiangyi commented on PR #10105:
URL: https://github.com/apache/paimon/pull/10105#issuecomment-5997688798

   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.
   


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