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]

Reply via email to