JingsongLi commented on code in PR #10133:
URL: https://github.com/apache/paimon/pull/10133#discussion_r4172611734


##########
paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java:
##########
@@ -85,6 +92,54 @@ public MetricGroup getMetricGroup() {
         return metricGroup;
     }
 
+    @VisibleForTesting
+    int activeCompactTimerCount() {
+        return compactTimers.size();
+    }
+
+    private Object compactTimerLock(long threadId) {
+        return compactTimerLocks.computeIfAbsent(threadId, ignored -> new 
Object());
+    }
+
+    private void releaseCompactTimer(long threadId) {
+        synchronized (compactTimerLock(threadId)) {
+            compactTimerRefCounts.compute(
+                    threadId,
+                    (id, count) -> {
+                        if (count == null || count <= 1) {
+                            compactTimers.remove(id);

Review Comment:
   **[P2] Preserve the busy window when a shared worker remains alive**
   
   Removing the timer when its last bucket reporter unregisters also changes 
the default `compaction.task-threads=1` behavior. An idle writer can be cleaned 
by `prepareCommit` shortly after its compaction finishes, while the shared 
executor remains alive. This deletion discards that worker's recent activity 
immediately instead of letting `compactionThreadBusy` decay over its existing 
60-second window.
   
   A deterministic timer probe recording 30 seconds of recent activity reports 
busy=50 before unregister. The PR head reports 0 immediately afterward; the 
base still reports 50. Rolling/sparse partitions can therefore under-report 
compaction utilization even without opting into parallelism.
   
   Please retain timers for the lifetime of SINGLE/FIXED_POOL workers, and 
limit retirement to workers that actually leave, such as dedicated PER_BUCKET 
executors. This also avoids introducing reporter-based timer reference 
bookkeeping into the unchanged default path.



##########
paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java:
##########
@@ -703,12 +713,94 @@ private void checkNumBuckets(String partInfo, int 
expected, int previous) {
         }
     }
 
-    private ExecutorService compactExecutor() {
+    private ExecutorService compactExecutor(BinaryRow partition, int bucket) {
+        if (externalCompactExecutor) {
+            return lazyCompactExecutor;
+        }
+
+        switch (compactionTaskExecutorMode) {
+            case PER_BUCKET:
+                return perBucketCompactExecutors.computeIfAbsent(
+                        new BucketCompactionExecutorKey(partition, bucket),
+                        key ->
+                                Executors.newSingleThreadExecutor(
+                                        new ExecutorThreadFactory(
+                                                
Thread.currentThread().getName()
+                                                        + "-compaction-bucket-"
+                                                        + key.bucket)));
+            case FIXED_POOL:
+                return sharedCompactionExecutor(compactionTaskThreads, 
"-compaction-pool");

Review Comment:
   **[P1] Isolate append compaction reader state before enabling cross-bucket 
concurrency**
   
   This routing also applies to `BucketedAppendFileStoreWrite`. Its bucket 
rewriters share `BaseAppendFileStoreWrite.readForCompact`, whose 
`RawFileSplitRead` caches one `FormatReaderMapping` per `(schemaId, format)`. 
The cached `CastFieldGetter[]` contains nested ROW/ARRAY/MAP casts that capture 
mutable `CastedRow`/`CastedArray`/`CastedMap` instances. Each 
`DataFileRecordReader` gets a separate outer row wrapper but shares these 
nested cast instances. Another bucket can therefore replace a nested wrapper's 
backing value while the first worker serializes it, silently writing the other 
bucket's data into the compacted file.
   
   I reproduced this on JDK 8 with an actual two-bucket append Parquet table: 
`(id=10, tags=[10])` and `(id=20, tags=[20])`, followed by the supported schema 
change `tags.element: INT -> BIGINT`. A latch-controlled interleaving around 
the real cached caster produced output files containing `(10, [20])` and `(20, 
[20])` with both `compaction.task-threads=2` and `-1`; the single-thread 
control preserved `(10, [10])` and `(20, [20])`.
   
   Please give each bucket/worker its own compaction reader and mutable cast 
state, or construct independent nested casts for each actual file reader. 
Making the cache a `ConcurrentHashMap` alone does not fix the shared mutable 
wrappers. Add a regression that checks rewritten nested values across buckets 
after schema evolution; the executor-routing tests do not exercise this path.



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