JingsongLi commented on PR #10133:
URL: https://github.com/apache/paimon/pull/10133#issuecomment-5952425655

   Reviewed head 4344ed7694. The feature has clear production value, and the 
previous generated-doc and timer-retirement feedback is addressed. Normal JDK 8 
verification passed 42 focused/adjacent tests, including worker churn, 
shared-worker handoff, executor ownership, and bucket serialization. I also ran 
actual write/update/delete -> compaction -> commit -> read workflows for all 
three modes with and without deletion vectors: observed concurrency was 1 / 2 / 
4 respectively, each bucket stayed serialized, all six workflows returned the 
correct 16 rows, and rewriters were closed. Generated option documentation, 
whitespace checks, and the exact-head CI are green.
   
   **[P2] Make shared compaction counters safe for the new worker concurrency** 
(`AbstractFileStoreWrite.java:800–802`, also the PER_BUCKET path). Multiple 
workers now call 
`CompactionMetrics.ReporterImpl.increaseCompactionsCompletedCount()` 
concurrently. In Flink, `FlinkMetricGroup.counter()` delegates to the default 
Flink `SimpleCounter`, whose `long` increment/decrement is not atomic. The core 
`TestMetricRegistry` uses an AtomicLong counter, so the new concurrency tests 
hide this production difference.
   
   I reproduced this through the real `FlinkMetricRegistry`/`FlinkMetricGroup` 
adapters and actual `CompactTask.call()` updates on Flink 1.20.4. One worker 
completed 10,000 tasks and reported exactly 10,000 / queued 0. Four workers 
completed 40,000 tasks, but three runs reported completion counts 39,952, 
39,986, and 39,913, with queued counts 5,607, 1,972, and 1,461 after all tasks 
finished. The new multi-worker completion-count loss is the introduced 
regression; the queued counter already had a main-thread/worker race, which 
this change also exposes under greater concurrency. Flink 2.2's default counter 
has the same implementation.
   
   Please synchronize updates to these shared counters or register a 
thread-safe counter through the Flink adapter, and add a regression using the 
actual production metric backend. This corrupts the completion/queue signals 
operators are explicitly told to use when sizing the new pool; it did not cause 
data loss or compaction failure in my tests.
   


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