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]