JingsongLi commented on PR #10105: URL: https://github.com/apache/paimon/pull/10105#issuecomment-5952813491
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. -- 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]
