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]