diff --git a/docs/docs/maintenance/metrics.md b/docs/docs/maintenance/metrics.md
index efd78e520d16..c2f18f0a8d79 100644
--- a/docs/docs/maintenance/metrics.md
+++ b/docs/docs/maintenance/metrics.md
@@ -444,7 +444,7 @@ Lookup metrics are available for local partial lookup. They are reported at look
| compactionThreadBusy |
Gauge |
- The maximum business of compaction threads in this task. Currently, there is only one compaction thread in each parallelism, so value of business ranges from 0 (idle) to 100 (compaction running all the time). |
+ The maximum busyness of compaction threads in this task, ranging from 0 (idle) to 100 when a single compaction thread is busy all the time. When compaction.task-threads is greater than 1 or set to -1 (per-bucket executors), multiple workers may be busy concurrently, so this value can exceed 100. |
| avgCompactionTime |
diff --git a/docs/docs/primary-key-table/compaction.md b/docs/docs/primary-key-table/compaction.md
index f69c587aa290..ae0a97360e65 100644
--- a/docs/docs/primary-key-table/compaction.md
+++ b/docs/docs/primary-key-table/compaction.md
@@ -70,6 +70,29 @@ publishes it. MOW batch readers can opt into merging pending data with
For `changelog-producer = lookup`, generated changelogs are also delayed. A compactor that cannot
keep up with sustained input will keep falling behind; relaxing waits does not add capacity.
+## Multi-thread async compaction
+
+When many buckets are assigned to the same write task (for example, one Flink sink subtask), the
+default async compaction uses **one shared thread** for all buckets. Under high write throughput,
+compaction may fall behind and level-0 files accumulate. This is especially visible in
+[MOW / deletion vectors mode](./table-mode#merge-on-write), where level-0 data becomes readable
+only after compaction publishes it.
+
+You can increase cross-bucket compaction parallelism with `compaction.task-threads`:
+
+- `1` (default): unchanged — one compaction thread per write task.
+- `N` (`N > 1`): a fixed thread pool of `N` threads shared by all buckets in the task.
+- `-1`: one dedicated compaction thread per active `(partition, bucket)` writer (highest
+ parallelism and memory use).
+
+Compaction **within the same bucket is always serialized**. Values `0` and negative integers other
+than `-1` are rejected.
+
+**Trade-offs:** more compaction threads increase TaskManager memory and I/O concurrency. Size
+TaskManager memory accordingly and monitor [compaction metrics](../maintenance/metrics#compaction-metrics)
+such as `avgLevel0FileCount`, `avgCompactionTime`, and `compactionQueuedCount`. Start with
+`N = 2` or `3` before using `-1`.
+
## Dedicated compaction job
Set `write-only = true` on ingest writers and run a
diff --git a/docs/docs/primary-key-table/table-mode.md b/docs/docs/primary-key-table/table-mode.md
index 78bb6eb24234..b18c7dabca28 100644
--- a/docs/docs/primary-key-table/table-mode.md
+++ b/docs/docs/primary-key-table/table-mode.md
@@ -108,6 +108,10 @@ By default, batch reads skip Level-0 files until lookup compaction publishes the
for this compaction by default. Asynchronous compaction or a dedicated compaction job can delay
visibility; see [Asynchronous Compaction](./compaction#asynchronous-compaction).
+When async compaction falls behind or data visibility latency is high, consider increasing
+`compaction.task-threads` to reduce visibility delay. See
+[Multi-thread async compaction](./compaction#multi-thread-async-compaction).
+
For batch scans, `deletion-vectors.merge-on-read = true` includes uncompacted data by merging it
at read time, with additional read cost. It does not change streaming changelog behavior.
diff --git a/docs/generated/core_configuration.html b/docs/generated/core_configuration.html
index a1c67478ac1d..389fee49308c 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -464,6 +464,12 @@
MemorySize |
When total size is smaller than this threshold, force a full compaction. |
+
+ compaction.task-threads |
+ 1 |
+ Integer |
+ Number of threads for async compaction in each write task (for example, each Flink sink subtask). 1 (default): all buckets in the task share one compaction thread. -1: one dedicated compaction thread per active (partition, bucket) writer in the task so different buckets compact in parallel. N (>1): a fixed thread pool of N threads shared by all buckets in the task. Compaction within the same bucket is still serialized. Values 0 and other negative integers except -1 are not allowed. Higher thread counts increase TaskManager memory pressure; monitor level-0 file count and compaction metrics. |
+
consumer-id |
(none) |
diff --git a/paimon-api/src/main/java/org/apache/paimon/CompactionTaskExecutorMode.java b/paimon-api/src/main/java/org/apache/paimon/CompactionTaskExecutorMode.java
new file mode 100644
index 000000000000..a4404bd16e85
--- /dev/null
+++ b/paimon-api/src/main/java/org/apache/paimon/CompactionTaskExecutorMode.java
@@ -0,0 +1,32 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon;
+
+/** Async compaction executor strategy for each write task (for example, a Flink sink subtask). */
+public enum CompactionTaskExecutorMode {
+
+ /** One shared compaction thread serializes compaction for all buckets in the task. */
+ SINGLE,
+
+ /** A fixed-size thread pool ({@code compaction.task-threads}) shared by all buckets. */
+ FIXED_POOL,
+
+ /** One dedicated compaction thread per active (partition, bucket) writer. */
+ PER_BUCKET
+}
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 98f0dfe048e6..3e376261181f 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -965,6 +965,22 @@ public InlineElement getDescription() {
.defaultValue(false)
.withDescription("Whether to force a compaction before commit.");
+ public static final ConfigOption COMPACTION_TASK_THREADS =
+ key("compaction.task-threads")
+ .intType()
+ .defaultValue(1)
+ .withDescription(
+ "Number of threads for async compaction in each write task (for example, "
+ + "each Flink sink subtask). "
+ + "1 (default): all buckets in the task share one compaction thread. "
+ + "-1: one dedicated compaction thread per active (partition, bucket) "
+ + "writer in the task so different buckets compact in parallel. "
+ + "N (>1): a fixed thread pool of N threads shared by all buckets in the task. "
+ + "Compaction within the same bucket is still serialized. "
+ + "Values 0 and other negative integers except -1 are not allowed. "
+ + "Higher thread counts increase TaskManager memory pressure; "
+ + "monitor level-0 file count and compaction metrics.");
+
public static final ConfigOption WRITE_SEQUENCE_NUMBER_INIT_MODE =
key("write.sequence-number-init-mode")
.enumType(SequenceNumberInitMode.class)
@@ -3896,6 +3912,26 @@ public boolean commitForceCompact() {
return options.get(COMMIT_FORCE_COMPACT);
}
+ public CompactionTaskExecutorMode compactionTaskExecutorMode() {
+ int threads = compactionTaskThreads();
+ if (threads == -1) {
+ return CompactionTaskExecutorMode.PER_BUCKET;
+ }
+ if (threads == 1) {
+ return CompactionTaskExecutorMode.SINGLE;
+ }
+ return CompactionTaskExecutorMode.FIXED_POOL;
+ }
+
+ public int compactionTaskThreads() {
+ int threads = options.get(COMPACTION_TASK_THREADS);
+ checkArgument(
+ threads == -1 || threads > 0,
+ "The option %s must be -1, 1, or any integer greater than 1.",
+ COMPACTION_TASK_THREADS.key());
+ return threads;
+ }
+
public SequenceNumberInitMode writeSequenceNumberInitMode() {
return options.get(WRITE_SEQUENCE_NUMBER_INIT_MODE);
}
diff --git a/paimon-api/src/test/java/org/apache/paimon/CoreOptionsCompactionTaskThreadsTest.java b/paimon-api/src/test/java/org/apache/paimon/CoreOptionsCompactionTaskThreadsTest.java
new file mode 100644
index 000000000000..ea9c48955c05
--- /dev/null
+++ b/paimon-api/src/test/java/org/apache/paimon/CoreOptionsCompactionTaskThreadsTest.java
@@ -0,0 +1,73 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon;
+
+import org.apache.paimon.options.Options;
+
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Tests for {@link CoreOptions#COMPACTION_TASK_THREADS}. */
+class CoreOptionsCompactionTaskThreadsTest {
+
+ @Test
+ void testDefaultIsSingleThreadMode() {
+ CoreOptions options = new CoreOptions(new Options());
+ assertThat(options.compactionTaskExecutorMode())
+ .isEqualTo(CompactionTaskExecutorMode.SINGLE);
+ assertThat(options.compactionTaskThreads()).isEqualTo(1);
+ }
+
+ @Test
+ void testFixedPoolMode() {
+ Options options = new Options();
+ options.set(CoreOptions.COMPACTION_TASK_THREADS, 3);
+ CoreOptions coreOptions = new CoreOptions(options);
+ assertThat(coreOptions.compactionTaskExecutorMode())
+ .isEqualTo(CompactionTaskExecutorMode.FIXED_POOL);
+ assertThat(coreOptions.compactionTaskThreads()).isEqualTo(3);
+ }
+
+ @Test
+ void testPerBucketMode() {
+ Options options = new Options();
+ options.set(CoreOptions.COMPACTION_TASK_THREADS, -1);
+ CoreOptions coreOptions = new CoreOptions(options);
+ assertThat(coreOptions.compactionTaskExecutorMode())
+ .isEqualTo(CompactionTaskExecutorMode.PER_BUCKET);
+ }
+
+ @Test
+ void testRejectZeroAndOtherNegativeValues() {
+ assertInvalid(0);
+ assertInvalid(-2);
+ assertInvalid(-100);
+ }
+
+ private static void assertInvalid(int threads) {
+ Options options = new Options();
+ options.set(CoreOptions.COMPACTION_TASK_THREADS, threads);
+ CoreOptions coreOptions = new CoreOptions(options);
+ assertThatThrownBy(coreOptions::compactionTaskExecutorMode)
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining(CoreOptions.COMPACTION_TASK_THREADS.key());
+ }
+}
diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java
index 0c0738c6173e..3af1cbbc0f16 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java
@@ -18,6 +18,7 @@
package org.apache.paimon.operation;
+import org.apache.paimon.CompactionTaskExecutorMode;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.KeyValue;
import org.apache.paimon.Snapshot;
@@ -61,7 +62,9 @@
import java.util.Iterator;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.OptionalLong;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.function.Function;
@@ -101,6 +104,11 @@ public abstract class AbstractFileStoreWrite implements FileStoreWrite {
protected final Map>> writers;
protected WriteRestore restore;
+ private final CompactionTaskExecutorMode compactionTaskExecutorMode;
+ private final int compactionTaskThreads;
+ private final Map perBucketCompactExecutors =
+ new ConcurrentHashMap<>();
+ private boolean externalCompactExecutor = false;
private ExecutorService lazyCompactExecutor;
private ExecutorService lazyPrimaryKeyIndexExecutor;
private boolean closeCompactExecutorWhenLeaving = true;
@@ -142,6 +150,8 @@ protected AbstractFileStoreWrite(
this.writerNumberMax = options.writeMaxWritersToSpill();
this.legacyPartitionName = options.legacyPartitionName();
this.options = options;
+ this.compactionTaskExecutorMode = options.compactionTaskExecutorMode();
+ this.compactionTaskThreads = options.compactionTaskThreads();
this.partitionTimestampValidator =
PartitionTimestampValidator.create(options, partitionType);
}
@@ -182,6 +192,7 @@ public void withIgnoreNumBucketCheck(boolean ignoreNumBucketCheck) {
public void withCompactExecutor(ExecutorService compactExecutor) {
this.lazyCompactExecutor = compactExecutor;
this.closeCompactExecutorWhenLeaving = false;
+ this.externalCompactExecutor = true;
}
@Override
@@ -309,6 +320,7 @@ public List prepareCommit(boolean waitCompaction, long commitIden
writerContainer.lastModifiedCommitIdentifier,
commitIdentifier);
}
+ releaseCompactionExecutor(partition, bucket);
writerContainer.writer.close();
if (writerContainer.primaryKeyIndexMaintainer != null) {
writerContainer.primaryKeyIndexMaintainer.close();
@@ -384,9 +396,7 @@ public void close() throws Exception {
// also left both thread pools running for the life of the process. None of the
// calls below throws, so the writer failure is never replaced by one of them.
writers.clear();
- if (lazyCompactExecutor != null && closeCompactExecutorWhenLeaving) {
- lazyCompactExecutor.shutdownNow();
- }
+ shutdownCompactionExecutors();
if (lazyPrimaryKeyIndexExecutor != null) {
lazyPrimaryKeyIndexExecutor.shutdownNow();
}
@@ -458,7 +468,7 @@ public void restore(List> states) {
state.dataFiles,
state.maxSequenceNumber,
state.commitIncrement,
- compactExecutor(),
+ compactExecutor(state.partition, state.bucket),
state.deletionVectorsMaintainer,
// Restore reconstructs writer state from checkpointed files, so do
// not ignore them.
@@ -591,7 +601,7 @@ private WriterContainer createWriterContainer(
startingMaxSequenceNumber(
getMaxSequenceNumber(restoreFiles), latestSnapshot),
null,
- compactExecutor(),
+ compactExecutor(partition, bucket),
dvMaintainer,
actualIgnorePreviousFiles);
notifyNewWriter(writer);
@@ -703,12 +713,98 @@ 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");
+ case SINGLE:
+ default:
+ return sharedCompactionExecutor(1, "-compaction");
+ }
+ }
+
+ private void releaseCompactionExecutor(BinaryRow partition, int bucket) {
+ if (compactionTaskExecutorMode != CompactionTaskExecutorMode.PER_BUCKET
+ || externalCompactExecutor) {
+ return;
+ }
+
+ if (compactionMetrics != null) {
+ compactionMetrics.retireCompactTimersForBucket(partition, bucket);
+ }
+
+ ExecutorService removed =
+ perBucketCompactExecutors.remove(
+ new BucketCompactionExecutorKey(partition, bucket));
+ if (removed != null) {
+ removed.shutdownNow();
+ }
+ }
+
+ private void shutdownCompactionExecutors() {
+ for (ExecutorService executor : perBucketCompactExecutors.values()) {
+ executor.shutdownNow();
+ }
+ perBucketCompactExecutors.clear();
+
+ if (lazyCompactExecutor != null && closeCompactExecutorWhenLeaving) {
+ lazyCompactExecutor.shutdownNow();
+ lazyCompactExecutor = null;
+ }
+ }
+
+ private static final class BucketCompactionExecutorKey {
+ private final BinaryRow partition;
+ private final int bucket;
+
+ private BucketCompactionExecutorKey(BinaryRow partition, int bucket) {
+ this.partition = partition;
+ this.bucket = bucket;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (o == null || getClass() != o.getClass()) {
+ return false;
+ }
+ BucketCompactionExecutorKey that = (BucketCompactionExecutorKey) o;
+ return bucket == that.bucket && Objects.equals(partition, that.partition);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(partition, bucket);
+ }
+ }
+
+ private ExecutorService sharedCompactionExecutor(int threads, String nameSuffix) {
if (lazyCompactExecutor == null) {
- lazyCompactExecutor =
- Executors.newSingleThreadScheduledExecutor(
- new ExecutorThreadFactory(
- Thread.currentThread().getName() + "-compaction"));
+ String threadNamePrefix = Thread.currentThread().getName() + nameSuffix;
+ if (threads <= 1) {
+ lazyCompactExecutor =
+ Executors.newSingleThreadScheduledExecutor(
+ new ExecutorThreadFactory(threadNamePrefix));
+ } else {
+ lazyCompactExecutor =
+ Executors.newFixedThreadPool(
+ threads, new ExecutorThreadFactory(threadNamePrefix));
+ }
}
return lazyCompactExecutor;
}
@@ -723,6 +819,21 @@ private ExecutorService primaryKeyIndexExecutor() {
return lazyPrimaryKeyIndexExecutor;
}
+ @VisibleForTesting
+ public ExecutorService compactExecutorForTesting(BinaryRow partition, int bucket) {
+ return compactExecutor(partition, bucket);
+ }
+
+ @VisibleForTesting
+ public int activePerBucketExecutorCountForTesting() {
+ return perBucketCompactExecutors.size();
+ }
+
+ @VisibleForTesting
+ public void releaseCompactionExecutorForTesting(BinaryRow partition, int bucket) {
+ releaseCompactionExecutor(partition, bucket);
+ }
+
@VisibleForTesting
public ExecutorService getCompactExecutor() {
return lazyCompactExecutor;
diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
index 07d15637fe95..d907c80bd801 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
@@ -19,6 +19,7 @@
package org.apache.paimon.operation;
import org.apache.paimon.AppendOnlyFileStore;
+import org.apache.paimon.CompactionTaskExecutorMode;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.append.AppendOnlyWriter;
import org.apache.paimon.append.cluster.Sorter;
@@ -83,6 +84,7 @@ public abstract class BaseAppendFileStoreWrite extends MemoryFileStoreWrite readForCompactByThread;
private final long schemaId;
private final FileFormat fileFormat;
private final FileStorePathFactory pathFactory;
@@ -120,6 +122,13 @@ public BaseAppendFileStoreWrite(
tableName);
this.fileIO = fileIO;
this.readForCompact = readForCompact;
+ if (options.compactionTaskExecutorMode() == CompactionTaskExecutorMode.SINGLE) {
+ this.readForCompactByThread = null;
+ } else {
+ this.readForCompactByThread =
+ ThreadLocal.withInitial(
+ () -> readForCompact.copyWithFreshReaderMappings(options));
+ }
this.schemaId = schemaId;
this.rowType = rowType;
this.writeType = rowType;
@@ -257,11 +266,24 @@ private SimpleColStatsCollector.Factory[] statsCollectors() {
@Override
public void close() throws Exception {
- if (blobFetchMetrics == null) {
- super.close();
- } else {
- IOUtils.closeAll(super::close, blobFetchMetrics::close);
+ try {
+ if (blobFetchMetrics == null) {
+ super.close();
+ } else {
+ IOUtils.closeAll(super::close, blobFetchMetrics::close);
+ }
+ } finally {
+ if (readForCompactByThread != null) {
+ readForCompactByThread.remove();
+ }
+ }
+ }
+
+ private RawFileSplitRead readForCompactReader() {
+ if (readForCompactByThread != null) {
+ return readForCompactByThread.get();
}
+ return readForCompact;
}
protected abstract CompactManager getCompactManager(
@@ -387,7 +409,7 @@ private RecordReaderIterator createFilesIterator(
@Nullable Map> dvFactories)
throws IOException {
return new RecordReaderIterator<>(
- readForCompact.createReader(partition, bucket, files, dvFactories));
+ readForCompactReader().createReader(partition, bucket, files, dvFactories));
}
@Override
diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java b/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java
index 84b32eaa1d97..2963246e759a 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java
@@ -117,6 +117,27 @@ public RawFileSplitRead(
this.readRowType = rowType;
}
+ /**
+ * Returns a reader with the same scan configuration but an empty format-reader mapping cache.
+ * Used so concurrent compaction workers do not share mutable cast state in cached mappings.
+ */
+ RawFileSplitRead copyWithFreshReaderMappings(CoreOptions coreOptions) {
+ RawFileSplitRead copy =
+ new RawFileSplitRead(
+ fileIO,
+ schemaManager,
+ schema,
+ readRowType,
+ formatDiscover,
+ pathFactory,
+ coreOptions);
+ copy.filters = filters;
+ copy.topN = topN;
+ copy.limit = limit;
+ copy.readBatchSizer = readBatchSizer;
+ return copy;
+ }
+
@Override
public SplitRead forceKeepDelete() {
return this;
diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java b/paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java
index ed6122b2b6dc..05c52c2bd274 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/metrics/CompactionMetrics.java
@@ -25,9 +25,11 @@
import org.apache.paimon.metrics.MetricRegistry;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.Map;
import java.util.Objects;
import java.util.Queue;
+import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.stream.DoubleStream;
@@ -66,25 +68,79 @@ public class CompactionMetrics {
private final MetricGroup metricGroup;
private final Map reporters;
private final Map compactTimers;
+ private final Map compactTimerLocks;
private final Queue compactionTimes;
private Counter compactionsCompletedCounter;
private Counter compactionsTotalCounter;
private Counter compactionsQueuedCounter;
+ private final Object sharedCounterLock = new Object();
public CompactionMetrics(MetricRegistry registry, String tableName) {
this.metricGroup = registry.createTableMetricGroup(GROUP_NAME, tableName);
this.reporters = new HashMap<>();
this.compactTimers = new ConcurrentHashMap<>();
+ this.compactTimerLocks = new ConcurrentHashMap<>();
this.compactionTimes = new ConcurrentLinkedQueue<>();
registerGenericCompactionMetrics();
}
+ /**
+ * Retire compact timers when an internally owned per-bucket compaction executor is released.
+ */
+ public void retireCompactTimersForBucket(BinaryRow partition, int bucket) {
+ PartitionAndBucket key = new PartitionAndBucket(partition, bucket);
+ ReporterImpl reporter = reporters.get(key);
+ if (reporter != null) {
+ reporter.retireCompactTimers();
+ }
+ }
+
@VisibleForTesting
public MetricGroup getMetricGroup() {
return metricGroup;
}
+ @VisibleForTesting
+ public 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)) {
+ compactTimers.remove(threadId);
+ compactTimerLocks.remove(threadId);
+ }
+ }
+
+ private void incrementCompactionsCompletedCount() {
+ synchronized (sharedCounterLock) {
+ compactionsCompletedCounter.inc();
+ }
+ }
+
+ private void incrementCompactionsTotalCount() {
+ synchronized (sharedCounterLock) {
+ compactionsTotalCounter.inc();
+ }
+ }
+
+ private void incrementCompactionsQueuedCount() {
+ synchronized (sharedCounterLock) {
+ compactionsQueuedCounter.inc();
+ }
+ }
+
+ private void decrementCompactionsQueuedCount() {
+ synchronized (sharedCounterLock) {
+ compactionsQueuedCounter.dec();
+ }
+ }
+
private void registerGenericCompactionMetrics() {
metricGroup.gauge(MAX_LEVEL0_FILE_COUNT, () -> getLevel0FileCountStream().max().orElse(-1));
metricGroup.gauge(
@@ -218,6 +274,7 @@ private class ReporterImpl implements Reporter {
private long totalFileCount = 0;
private long sortBufferUsedBytes = 0;
private double sortBufferUtilisationPercent = 0.0;
+ private final Set compactThreadIds = new HashSet<>();
private ReporterImpl(PartitionAndBucket key) {
this.key = key;
@@ -226,9 +283,12 @@ private ReporterImpl(PartitionAndBucket key) {
@Override
public CompactTimer getCompactTimer() {
- return compactTimers.computeIfAbsent(
- Thread.currentThread().getId(),
- ignore -> new CompactTimer(BUSY_MEASURE_MILLIS));
+ long threadId = Thread.currentThread().getId();
+ synchronized (compactTimerLock(threadId)) {
+ compactThreadIds.add(threadId);
+ return compactTimers.computeIfAbsent(
+ threadId, ignore -> new CompactTimer(BUSY_MEASURE_MILLIS));
+ }
}
@Override
@@ -275,26 +335,34 @@ public void reportSortBufferMetrics(long usedBytes, long totalBytes) {
@Override
public void increaseCompactionsCompletedCount() {
- compactionsCompletedCounter.inc();
+ CompactionMetrics.this.incrementCompactionsCompletedCount();
}
@Override
public void increaseCompactionsTotalCount() {
- compactionsTotalCounter.inc();
+ CompactionMetrics.this.incrementCompactionsTotalCount();
}
@Override
public void increaseCompactionsQueuedCount() {
- compactionsQueuedCounter.inc();
+ CompactionMetrics.this.incrementCompactionsQueuedCount();
}
@Override
public void decreaseCompactionsQueuedCount() {
- compactionsQueuedCounter.dec();
+ CompactionMetrics.this.decrementCompactionsQueuedCount();
+ }
+
+ private void retireCompactTimers() {
+ for (Long threadId : compactThreadIds) {
+ releaseCompactTimer(threadId);
+ }
+ compactThreadIds.clear();
}
@Override
public void unregister() {
+ compactThreadIds.clear();
reporters.remove(key);
}
}
diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerBucketSerializationTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerBucketSerializationTest.java
new file mode 100644
index 000000000000..491a3f905fd6
--- /dev/null
+++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerBucketSerializationTest.java
@@ -0,0 +1,125 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.mergetree.compact;
+
+import org.apache.paimon.compact.CompactResult;
+import org.apache.paimon.compact.CompactUnit;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.io.DataFileTestUtils;
+import org.apache.paimon.mergetree.LevelSortedRun;
+import org.apache.paimon.mergetree.Levels;
+import org.apache.paimon.mergetree.SortedRun;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/** Tests that a single {@link MergeTreeCompactManager} serializes compaction per bucket. */
+class MergeTreeCompactManagerBucketSerializationTest {
+
+ private ExecutorService executorService;
+
+ @AfterEach
+ void tearDown() {
+ if (executorService != null) {
+ executorService.shutdownNow();
+ executorService = null;
+ }
+ }
+
+ @Test
+ void testSecondTriggerIgnoredWhileCompactionRunning() throws Exception {
+ executorService = Executors.newSingleThreadExecutor();
+ AtomicInteger rewriteCalls = new AtomicInteger();
+ CountDownLatch rewriteStarted = new CountDownLatch(1);
+ CountDownLatch unblockRewrite = new CountDownLatch(1);
+
+ DataFileMeta file1 = DataFileTestUtils.newFile(0, 1, 3, 3L);
+ DataFileMeta file2 = DataFileTestUtils.newFile(0, 4, 6, 6L);
+ LevelSortedRun run1 = new LevelSortedRun(0, SortedRun.fromSingle(file1));
+ LevelSortedRun run2 = new LevelSortedRun(0, SortedRun.fromSingle(file2));
+ List runs = Arrays.asList(run1, run2);
+
+ Levels levels = mock(Levels.class);
+ when(levels.levelSortedRuns()).thenReturn(runs);
+ when(levels.numberOfLevels()).thenReturn(3);
+ when(levels.nonEmptyHighestLevel()).thenReturn(0);
+
+ CompactStrategy strategy = mock(CompactStrategy.class);
+ when(strategy.pick(anyInt(), any()))
+ .thenReturn(Optional.of(CompactUnit.fromLevelRuns(1, runs)));
+
+ CompactRewriter rewriter = mock(CompactRewriter.class);
+ when(rewriter.rewrite(anyInt(), anyBoolean(), any()))
+ .thenAnswer(
+ invocation -> {
+ rewriteCalls.incrementAndGet();
+ rewriteStarted.countDown();
+ assertThat(unblockRewrite.await(30, TimeUnit.SECONDS)).isTrue();
+ return new CompactResult(
+ Collections.emptyList(), Collections.emptyList());
+ });
+
+ MergeTreeCompactManager manager =
+ new MergeTreeCompactManager(
+ executorService,
+ levels,
+ strategy,
+ Comparator.comparingInt(row -> row.getInt(0)),
+ 1024 * 1024,
+ 5,
+ rewriter,
+ null,
+ null,
+ false,
+ false,
+ null,
+ false,
+ false,
+ "");
+
+ manager.triggerCompaction(false);
+ assertThat(rewriteStarted.await(30, TimeUnit.SECONDS)).isTrue();
+ assertThat(manager.compactNotCompleted()).isTrue();
+
+ manager.triggerCompaction(false);
+ assertThat(rewriteCalls.get()).isEqualTo(1);
+
+ unblockRewrite.countDown();
+ manager.getCompactionResult(true);
+ assertThat(manager.compactNotCompleted()).isFalse();
+ }
+}
diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/CompactionTaskExecutorRoutingTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/CompactionTaskExecutorRoutingTest.java
new file mode 100644
index 000000000000..709516fb8bee
--- /dev/null
+++ b/paimon-core/src/test/java/org/apache/paimon/operation/CompactionTaskExecutorRoutingTest.java
@@ -0,0 +1,239 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.operation;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.deletionvectors.BucketedDvMaintainer;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.metrics.TestMetricRegistry;
+import org.apache.paimon.operation.metrics.CompactionMetrics;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.CommitIncrement;
+import org.apache.paimon.utils.RecordWriter;
+import org.apache.paimon.utils.SnapshotManager;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import javax.annotation.Nullable;
+
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Function;
+
+import static org.apache.paimon.data.BinaryRow.EMPTY_ROW;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.mock;
+
+/** Focused behavioral tests for {@code compaction.task-threads} executor routing. */
+class CompactionTaskExecutorRoutingTest {
+
+ private ExecutorService externalExecutor;
+
+ @AfterEach
+ void tearDown() {
+ if (externalExecutor != null) {
+ externalExecutor.shutdownNow();
+ externalExecutor = null;
+ }
+ }
+
+ @Test
+ void testFixedPoolUsesSharedThreadPoolWithConfiguredSize() throws Exception {
+ ExecutorRoutingWrite write = new ExecutorRoutingWrite(coreOptions(2));
+ ExecutorService bucket0 = write.compactExecutorForTesting(EMPTY_ROW, 0);
+ ExecutorService bucket1 = write.compactExecutorForTesting(EMPTY_ROW, 1);
+ assertThat(bucket0).isSameAs(bucket1);
+ assertThat(bucket0).isInstanceOf(ThreadPoolExecutor.class);
+ assertThat(((ThreadPoolExecutor) bucket0).getMaximumPoolSize()).isEqualTo(2);
+ write.close();
+ }
+
+ @Test
+ void testFixedPoolAllowsCrossBucketParallelism() throws Exception {
+ ExecutorRoutingWrite write = new ExecutorRoutingWrite(coreOptions(2));
+ ExecutorService pool = write.compactExecutorForTesting(EMPTY_ROW, 0);
+
+ CountDownLatch firstStarted = new CountDownLatch(1);
+ CountDownLatch releaseFirst = new CountDownLatch(1);
+ CountDownLatch secondStarted = new CountDownLatch(1);
+
+ Future> first =
+ pool.submit(
+ () -> {
+ firstStarted.countDown();
+ try {
+ releaseFirst.await(30, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ });
+ Future> second =
+ pool.submit(
+ () -> {
+ secondStarted.countDown();
+ });
+
+ assertThat(firstStarted.await(30, TimeUnit.SECONDS)).isTrue();
+ assertThat(secondStarted.await(30, TimeUnit.SECONDS)).isTrue();
+
+ releaseFirst.countDown();
+ first.get(30, TimeUnit.SECONDS);
+ second.get(30, TimeUnit.SECONDS);
+ write.close();
+ }
+
+ @Test
+ void testPerBucketUsesDedicatedExecutorsAndReusesSameBucketExecutor() throws Exception {
+ BinaryRow partition = EMPTY_ROW.copy();
+ ExecutorRoutingWrite write = new ExecutorRoutingWrite(coreOptions(-1));
+
+ ExecutorService bucket0 = write.compactExecutorForTesting(partition, 0);
+ ExecutorService bucket1 = write.compactExecutorForTesting(partition, 1);
+ assertThat(bucket0).isNotSameAs(bucket1);
+ assertThat(write.activePerBucketExecutorCountForTesting()).isEqualTo(2);
+
+ ExecutorService bucket0Again = write.compactExecutorForTesting(partition, 0);
+ assertThat(bucket0Again).isSameAs(bucket0);
+ write.close();
+ }
+
+ @Test
+ void testPerBucketReleaseShutsDownAndRecreatesExecutor() throws Exception {
+ BinaryRow partition = EMPTY_ROW.copy();
+ ExecutorRoutingWrite write = new ExecutorRoutingWrite(coreOptions(-1));
+ ExecutorService first = write.compactExecutorForTesting(partition, 0);
+ assertThat(write.activePerBucketExecutorCountForTesting()).isEqualTo(1);
+
+ write.releaseCompactionExecutorForTesting(partition, 0);
+ assertThat(write.activePerBucketExecutorCountForTesting()).isZero();
+ assertThat(first.isShutdown()).isTrue();
+
+ ExecutorService second = write.compactExecutorForTesting(partition, 0);
+ assertThat(second).isNotSameAs(first);
+ write.close();
+ }
+
+ @Test
+ void testExternalExecutorIsSharedAndNotClosedByWrite() throws Exception {
+ externalExecutor = Executors.newSingleThreadExecutor();
+ ExecutorRoutingWrite write = new ExecutorRoutingWrite(coreOptions(-1));
+ write.withCompactExecutor(externalExecutor);
+
+ ExecutorService bucket0 = write.compactExecutorForTesting(EMPTY_ROW, 0);
+ ExecutorService bucket1 = write.compactExecutorForTesting(EMPTY_ROW, 1);
+ assertThat(bucket0).isSameAs(externalExecutor);
+ assertThat(bucket1).isSameAs(externalExecutor);
+ assertThat(write.activePerBucketExecutorCountForTesting()).isZero();
+
+ write.close();
+ assertThat(externalExecutor.isShutdown()).isFalse();
+ }
+
+ @Test
+ void testInternalPerBucketExecutorReleaseRetiresCompactTimer() throws Exception {
+ BinaryRow partition = EMPTY_ROW.copy();
+ ExecutorRoutingWrite write = new ExecutorRoutingWrite(coreOptions(-1));
+ write.withMetricRegistry(new TestMetricRegistry());
+
+ ExecutorService bucketExecutor = write.compactExecutorForTesting(partition, 0);
+ CompactionMetrics metrics = write.compactionMetrics();
+ CompactionMetrics.Reporter reporter = metrics.createReporter(partition, 0);
+ bucketExecutor
+ .submit(
+ () -> {
+ reporter.getCompactTimer().start();
+ reporter.getCompactTimer().finish();
+ })
+ .get(30, TimeUnit.SECONDS);
+
+ assertThat(metrics.activeCompactTimerCount()).isEqualTo(1);
+ write.releaseCompactionExecutorForTesting(partition, 0);
+ assertThat(metrics.activeCompactTimerCount()).isZero();
+ reporter.unregister();
+ write.close();
+ }
+
+ @Test
+ void testExternalExecutorDisablesCompactTimerRetirementOnPerBucketMode() throws Exception {
+ externalExecutor = Executors.newSingleThreadExecutor();
+ ExecutorRoutingWrite write = new ExecutorRoutingWrite(coreOptions(-1));
+ write.withMetricRegistry(new TestMetricRegistry());
+ write.withCompactExecutor(externalExecutor);
+
+ CompactionMetrics metrics = write.compactionMetrics();
+ CompactionMetrics.Reporter reporter = metrics.createReporter(EMPTY_ROW, 0);
+ externalExecutor
+ .submit(
+ () -> {
+ reporter.getCompactTimer().start();
+ reporter.getCompactTimer().finish();
+ })
+ .get(30, TimeUnit.SECONDS);
+ reporter.unregister();
+ assertThat(metrics.activeCompactTimerCount()).isEqualTo(1);
+ write.close();
+ }
+
+ private static CoreOptions coreOptions(int compactionTaskThreads) {
+ Options options = new Options();
+ options.set(CoreOptions.COMPACTION_TASK_THREADS, compactionTaskThreads);
+ return new CoreOptions(options);
+ }
+
+ private static class ExecutorRoutingWrite extends AbstractFileStoreWrite {
+
+ private ExecutorRoutingWrite(CoreOptions coreOptions) {
+ super(
+ mock(SnapshotManager.class),
+ mock(FileStoreScan.class),
+ null,
+ null,
+ null,
+ "test-table",
+ coreOptions,
+ RowType.of());
+ }
+
+ @Override
+ protected Function, Boolean> createWriterCleanChecker() {
+ return writer -> false;
+ }
+
+ @Override
+ protected RecordWriter createWriter(
+ BinaryRow partition,
+ int bucket,
+ List restoreFiles,
+ long restoredMaxSeqNumber,
+ @Nullable CommitIncrement restoreIncrement,
+ ExecutorService compactExecutor,
+ @Nullable BucketedDvMaintainer deletionVectorsMaintainer,
+ boolean ignorePreviousFiles) {
+ return mock(RecordWriter.class);
+ }
+ }
+}
diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/ParallelAppendCompactionReaderIsolationTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/ParallelAppendCompactionReaderIsolationTest.java
new file mode 100644
index 000000000000..bd663b032ddb
--- /dev/null
+++ b/paimon-core/src/test/java/org/apache/paimon/operation/ParallelAppendCompactionReaderIsolationTest.java
@@ -0,0 +1,192 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.operation;
+
+import org.apache.paimon.AppendOnlyFileStore;
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.catalog.Catalog;
+import org.apache.paimon.catalog.FileSystemCatalog;
+import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.BinaryRowWriter;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.deletionvectors.DeletionVector;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.reader.RecordReaderIterator;
+import org.apache.paimon.schema.FileSystemSchemaManager;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.SchemaChange;
+import org.apache.paimon.schema.SchemaManager;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.sink.StreamTableCommit;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.utils.IOExceptionSupplier;
+
+import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import javax.annotation.Nullable;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.stream.Collectors;
+
+import static org.apache.paimon.CoreOptions.BUCKET;
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * Regression for parallel append compaction: each worker must use isolated compaction readers so
+ * nested cast state is not shared across buckets.
+ */
+public class ParallelAppendCompactionReaderIsolationTest {
+
+ @TempDir java.nio.file.Path tempDir;
+
+ @ParameterizedTest
+ @ValueSource(ints = {2, -1})
+ public void testParallelCompactRewriteAfterNestedSchemaEvolution(int compactionTaskThreads)
+ throws Exception {
+ Path warehouse = new Path(tempDir.toString());
+ Catalog catalog = new FileSystemCatalog(LocalFileIO.create(), warehouse);
+ Identifier identifier = Identifier.create("default", "append_parallel_compact");
+ catalog.createDatabase("default", true);
+
+ Schema schema =
+ Schema.newBuilder()
+ .column("pt", DataTypes.INT())
+ .column("bk", DataTypes.INT())
+ .column("marker", DataTypes.INT())
+ .column(
+ "payload",
+ DataTypes.ROW(DataTypes.FIELD(0, "val", DataTypes.INT())))
+ .partitionKeys("pt")
+ .option(BUCKET.key(), "2")
+ .option("bucket-key", "bk")
+ .option("file.format", "parquet")
+ .option(CoreOptions.COMPACTION_TASK_THREADS.key(), "1")
+ .option("target-file-size", "256 b")
+ .option("write-buffer-size", "256 b")
+ .build();
+ catalog.createTable(identifier, schema, false);
+
+ FileStoreTable table = (FileStoreTable) catalog.getTable(identifier);
+ BinaryRow partition = partition(0);
+ String commitUser = "user-1";
+
+ try (StreamTableCommit commit = table.newStreamWriteBuilder().newCommit()) {
+ BaseAppendFileStoreWrite write =
+ (BaseAppendFileStoreWrite) table.store().newWrite(commitUser);
+ for (int i = 0; i < 24; i++) {
+ int bucket = i % 2;
+ int marker = bucket == 0 ? 100 : 200;
+ write.write(partition, bucket, GenericRow.of(0, bucket, marker, GenericRow.of(i)));
+ commit.commit(i, write.prepareCommit(false, i));
+ }
+ write.close();
+ }
+
+ Path tablePath = new Path(warehouse, "default.db/append_parallel_compact");
+ SchemaManager schemaManager = new FileSystemSchemaManager(LocalFileIO.create(), tablePath);
+ schemaManager.commitChanges(
+ SchemaChange.updateColumnType(
+ new String[] {"payload", "val"}, DataTypes.BIGINT(), false));
+
+ table =
+ (FileStoreTable)
+ catalog.getTable(identifier)
+ .copy(
+ Collections.singletonMap(
+ CoreOptions.COMPACTION_TASK_THREADS.key(),
+ String.valueOf(compactionTaskThreads)));
+
+ List bucket0Files = dataFiles(table, partition, 0);
+ List bucket1Files = dataFiles(table, partition, 1);
+ assertThat(bucket0Files).isNotEmpty();
+ assertThat(bucket1Files).isNotEmpty();
+
+ BaseAppendFileStoreWrite write =
+ (BaseAppendFileStoreWrite) table.store().newWrite(commitUser);
+ ExecutorService pool = Executors.newFixedThreadPool(2);
+ try {
+ Future> bucket0Future =
+ pool.submit(() -> write.compactRewrite(partition, 0, null, bucket0Files));
+ Future> bucket1Future =
+ pool.submit(() -> write.compactRewrite(partition, 1, null, bucket1Files));
+ assertRowsHaveMarker(table, partition, 0, bucket0Future.get(), 100);
+ assertRowsHaveMarker(table, partition, 1, bucket1Future.get(), 200);
+ } finally {
+ pool.shutdownNow();
+ write.close();
+ }
+ }
+
+ private static List dataFiles(
+ FileStoreTable table, BinaryRow partition, int bucket) {
+ List files = new ArrayList<>();
+ for (Split split : table.newScan().plan().splits()) {
+ DataSplit dataSplit = (DataSplit) split;
+ if (dataSplit.bucket() == bucket && dataSplit.partition().equals(partition)) {
+ files.addAll(dataSplit.dataFiles());
+ }
+ }
+ return files;
+ }
+
+ private static void assertRowsHaveMarker(
+ FileStoreTable table,
+ BinaryRow partition,
+ int bucket,
+ List files,
+ int expectedMarker)
+ throws Exception {
+ assertThat(files).isNotEmpty();
+ RawFileSplitRead read = ((AppendOnlyFileStore) table.store()).newRead();
+ @Nullable Map> dvFactories = null;
+ try (RecordReaderIterator iterator =
+ new RecordReaderIterator<>(
+ read.createReader(partition, bucket, files, dvFactories))) {
+ List markers = new ArrayList<>();
+ while (iterator.hasNext()) {
+ markers.add(iterator.next().getInt(2));
+ }
+ assertThat(markers).isNotEmpty();
+ assertThat(markers.stream().distinct().collect(Collectors.toList()))
+ .containsExactly(expectedMarker);
+ }
+ }
+
+ private static BinaryRow partition(int pt) {
+ BinaryRow row = new BinaryRow(1);
+ BinaryRowWriter writer = new BinaryRowWriter(row);
+ writer.writeInt(0, pt);
+ writer.complete();
+ return row;
+ }
+}
diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/metrics/CompactionMetricsTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/metrics/CompactionMetricsTest.java
index 380eac4b918a..658da44ae446 100644
--- a/paimon-core/src/test/java/org/apache/paimon/operation/metrics/CompactionMetricsTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/operation/metrics/CompactionMetricsTest.java
@@ -28,7 +28,10 @@
import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.metrics.Counter;
import org.apache.paimon.metrics.Gauge;
+import org.apache.paimon.metrics.Histogram;
import org.apache.paimon.metrics.Metric;
+import org.apache.paimon.metrics.MetricGroup;
+import org.apache.paimon.metrics.MetricRegistry;
import org.apache.paimon.metrics.TestMetricRegistry;
import org.apache.paimon.operation.AbstractFileStoreWrite;
import org.apache.paimon.options.Options;
@@ -47,8 +50,16 @@
import java.util.Arrays;
import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
import java.util.UUID;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
import java.util.concurrent.ThreadLocalRandom;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
import static org.assertj.core.api.Assertions.assertThat;
@@ -257,6 +268,227 @@ public void testTotalFileSizeForPrimaryKeyTables() throws Exception {
commit.close();
}
+ @Test
+ public void testCompactTimersRetiredAfterPerBucketWorkerChurn() throws Exception {
+ CompactionMetrics metrics = new CompactionMetrics(new TestMetricRegistry(), "myTable");
+ for (int i = 0; i < 32; i++) {
+ ExecutorService worker = Executors.newSingleThreadExecutor();
+ CompactionMetrics.Reporter reporter = metrics.createReporter(BinaryRow.EMPTY_ROW, i);
+ try {
+ worker.submit(
+ () -> {
+ reporter.getCompactTimer().start();
+ reporter.getCompactTimer().finish();
+ })
+ .get(30, TimeUnit.SECONDS);
+ } finally {
+ metrics.retireCompactTimersForBucket(BinaryRow.EMPTY_ROW, i);
+ reporter.unregister();
+ worker.shutdownNow();
+ }
+ }
+ assertThat(metrics.activeCompactTimerCount()).isZero();
+ }
+
+ @Test
+ public void testCompactTimerConcurrentUnregisterAndStartOnSharedWorker() throws Exception {
+ for (int attempt = 0; attempt < 200; attempt++) {
+ CompactionMetrics metrics = new CompactionMetrics(new TestMetricRegistry(), "myTable");
+ ExecutorService worker = Executors.newSingleThreadExecutor();
+ CompactionMetrics.Reporter retiring = metrics.createReporter(BinaryRow.EMPTY_ROW, 0);
+ CompactionMetrics.Reporter starting = metrics.createReporter(BinaryRow.EMPTY_ROW, 1);
+ CountDownLatch compactionStarted = new CountDownLatch(1);
+ CountDownLatch allowWorkerContinue = new CountDownLatch(1);
+ AtomicReference workerError = new AtomicReference<>();
+
+ Future> compaction =
+ worker.submit(
+ () -> {
+ try {
+ retiring.getCompactTimer().start();
+ retiring.getCompactTimer().finish();
+ compactionStarted.countDown();
+ allowWorkerContinue.await(30, TimeUnit.SECONDS);
+ starting.getCompactTimer().start();
+ starting.getCompactTimer().finish();
+ } catch (Throwable t) {
+ workerError.set(t);
+ }
+ });
+
+ assertThat(compactionStarted.await(30, TimeUnit.SECONDS)).isTrue();
+ retiring.unregister();
+ allowWorkerContinue.countDown();
+
+ compaction.get(30, TimeUnit.SECONDS);
+ worker.shutdownNow();
+
+ assertThat(workerError.get()).isNull();
+ starting.unregister();
+ assertThat(metrics.activeCompactTimerCount()).isEqualTo(1);
+ }
+ }
+
+ @Test
+ public void testSharedCompactionCountersWithConcurrentNonAtomicBackend() throws Exception {
+ Map counters = new HashMap<>();
+ MetricRegistry registry =
+ (groupName, variables) ->
+ new MetricGroup() {
+ @Override
+ public Counter counter(String name) {
+ return counters.computeIfAbsent(
+ name, n -> new NonThreadSafeCounter());
+ }
+
+ @Override
+ public Gauge gauge(String name, Gauge gauge) {
+ return gauge;
+ }
+
+ @Override
+ public Histogram histogram(String name, int windowSize) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public Map getAllVariables() {
+ return variables;
+ }
+
+ @Override
+ public String getGroupName() {
+ return groupName;
+ }
+
+ @Override
+ public Map getMetrics() {
+ return Collections.emptyMap();
+ }
+
+ @Override
+ public void close() {}
+ };
+
+ CompactionMetrics metrics = new CompactionMetrics(registry, "myTable");
+ CompactionMetrics.Reporter[] reporters = new CompactionMetrics.Reporter[4];
+ for (int i = 0; i < reporters.length; i++) {
+ reporters[i] = metrics.createReporter(BinaryRow.EMPTY_ROW, i);
+ }
+
+ ExecutorService workers = Executors.newFixedThreadPool(4);
+ try {
+ Future>[] futures = new Future>[reporters.length];
+ for (int i = 0; i < reporters.length; i++) {
+ CompactionMetrics.Reporter reporter = reporters[i];
+ futures[i] =
+ workers.submit(
+ () -> {
+ for (int j = 0; j < 10_000; j++) {
+ reporter.increaseCompactionsCompletedCount();
+ reporter.increaseCompactionsTotalCount();
+ reporter.increaseCompactionsQueuedCount();
+ reporter.decreaseCompactionsQueuedCount();
+ }
+ });
+ }
+ for (Future> future : futures) {
+ future.get(60, TimeUnit.SECONDS);
+ }
+ } finally {
+ workers.shutdownNow();
+ }
+
+ assertThat(counters.get(CompactionMetrics.COMPACTION_COMPLETED_COUNT).getCount())
+ .isEqualTo(40_000L);
+ assertThat(counters.get(CompactionMetrics.COMPACTION_TOTAL_COUNT).getCount())
+ .isEqualTo(40_000L);
+ assertThat(counters.get(CompactionMetrics.COMPACTION_QUEUED_COUNT).getCount())
+ .isEqualTo(0L);
+ }
+
+ @Test
+ public void testCompactTimerNotRetiredOnReporterUnregisterForSharedExecutor() throws Exception {
+ CompactionMetrics metrics = new CompactionMetrics(new TestMetricRegistry(), "myTable");
+ ExecutorService sharedPool = Executors.newSingleThreadExecutor();
+ CompactionMetrics.Reporter first = metrics.createReporter(BinaryRow.EMPTY_ROW, 0);
+ CompactionMetrics.Reporter second = metrics.createReporter(BinaryRow.EMPTY_ROW, 1);
+ try {
+ sharedPool
+ .submit(
+ () -> {
+ first.getCompactTimer().start();
+ first.getCompactTimer().finish();
+ second.getCompactTimer().start();
+ second.getCompactTimer().finish();
+ })
+ .get(30, TimeUnit.SECONDS);
+ first.unregister();
+ second.unregister();
+ assertThat(metrics.activeCompactTimerCount()).isEqualTo(1);
+ } finally {
+ sharedPool.shutdownNow();
+ }
+ }
+
+ @Test
+ public void testCompactTimerKeptWhileSharedCompactionThreadInUse() throws Exception {
+ CompactionMetrics metrics = new CompactionMetrics(new TestMetricRegistry(), "myTable");
+ ExecutorService sharedPool = Executors.newFixedThreadPool(1);
+ CompactionMetrics.Reporter first = metrics.createReporter(BinaryRow.EMPTY_ROW, 0);
+ CompactionMetrics.Reporter second = metrics.createReporter(BinaryRow.EMPTY_ROW, 1);
+ try {
+ sharedPool
+ .submit(
+ () -> {
+ first.getCompactTimer().start();
+ first.getCompactTimer().finish();
+ second.getCompactTimer().start();
+ second.getCompactTimer().finish();
+ })
+ .get(30, TimeUnit.SECONDS);
+ first.unregister();
+ assertThat(metrics.activeCompactTimerCount()).isEqualTo(1);
+ second.unregister();
+ assertThat(metrics.activeCompactTimerCount()).isEqualTo(1);
+ } finally {
+ sharedPool.shutdownNow();
+ }
+ }
+
+ /**
+ * Mimics Flink's {@code org.apache.flink.metrics.SimpleCounter}, which is not safe under
+ * concurrent updates from multiple compaction worker threads.
+ */
+ private static class NonThreadSafeCounter implements Counter {
+ private long count;
+
+ @Override
+ public void inc() {
+ count++;
+ }
+
+ @Override
+ public void inc(long n) {
+ count += n;
+ }
+
+ @Override
+ public void dec() {
+ count--;
+ }
+
+ @Override
+ public void dec(long n) {
+ count -= n;
+ }
+
+ @Override
+ public long getCount() {
+ return count;
+ }
+ }
+
private Object getMetric(CompactionMetrics metrics, String metricName) {
Metric metric = metrics.getMetricGroup().getMetrics().get(metricName);
if (metric instanceof Gauge) {