This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new c852c1a6ae [core] Optimize compaction duration log
c852c1a6ae is described below
commit c852c1a6ae5a7e235f79063f78c04c82baa7852d
Author: JingsongLi <[email protected]>
AuthorDate: Fri Jul 3 20:54:09 2026 +0800
[core] Optimize compaction duration log
---
.../org/apache/paimon/compact/CompactTask.java | 17 +++-----
.../mergetree/compact/FileRewriteCompactTask.java | 5 ++-
.../mergetree/compact/MergeTreeCompactManager.java | 20 ++++-----
.../compact/MergeTreeCompactManagerFactory.java | 48 ++++++++++------------
.../mergetree/compact/MergeTreeCompactTask.java | 5 ++-
.../paimon/operation/AbstractFileStoreWrite.java | 27 ++++++------
.../org/apache/paimon/operation/WriteRestore.java | 11 +----
.../apache/paimon/mergetree/MergeTreeTestBase.java | 6 ++-
.../compact/MergeTreeCompactManagerTest.java | 12 ++++--
.../sink/coordinator/TableWriteCoordinator.java | 20 +--------
10 files changed, 71 insertions(+), 100 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/compact/CompactTask.java
b/paimon-core/src/main/java/org/apache/paimon/compact/CompactTask.java
index faaf038106..229b25324b 100644
--- a/paimon-core/src/main/java/org/apache/paimon/compact/CompactTask.java
+++ b/paimon-core/src/main/java/org/apache/paimon/compact/CompactTask.java
@@ -36,16 +36,11 @@ public abstract class CompactTask implements
Callable<CompactResult> {
private static final Logger LOG =
LoggerFactory.getLogger(CompactTask.class);
@Nullable private final CompactionMetrics.Reporter metricsReporter;
+ private final String bucketInfo;
- private String logInfo = "";
-
- public CompactTask(@Nullable CompactionMetrics.Reporter metricsReporter) {
+ public CompactTask(@Nullable CompactionMetrics.Reporter metricsReporter,
String bucketInfo) {
this.metricsReporter = metricsReporter;
- }
-
- /** Set additional compact task information for logging purposes. */
- public void setLogInfo(String logInfo) {
- this.logInfo = logInfo;
+ this.bucketInfo = bucketInfo;
}
@Override
@@ -53,7 +48,7 @@ public abstract class CompactTask implements
Callable<CompactResult> {
MetricUtils.safeCall(this::startTimer, LOG);
LOG.info(
"Paimon compact task started: {}, taskType={}",
- logInfo,
+ bucketInfo,
getClass().getSimpleName());
try {
long startMillis = System.currentTimeMillis();
@@ -82,7 +77,7 @@ public abstract class CompactTask implements
Callable<CompactResult> {
LOG.info(
"Paimon compact task finished: {}, taskType={}, "
+ "inputFiles={}, inputBytes={}, outputFiles={},
outputBytes={}, durationMs={}",
- logInfo,
+ bucketInfo,
getClass().getSimpleName(),
result.before().size(),
result.before().stream().mapToLong(DataFileMeta::fileSize).sum(),
@@ -97,7 +92,7 @@ public abstract class CompactTask implements
Callable<CompactResult> {
} catch (Exception e) {
LOG.warn(
"Paimon compact task failed: {}, taskType={}",
- logInfo,
+ bucketInfo,
getClass().getSimpleName(),
e);
throw e;
diff --git
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FileRewriteCompactTask.java
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FileRewriteCompactTask.java
index 620ea0748a..686f640d49 100644
---
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FileRewriteCompactTask.java
+++
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/FileRewriteCompactTask.java
@@ -38,17 +38,20 @@ public class FileRewriteCompactTask extends CompactTask {
private final int outputLevel;
private final List<DataFileMeta> files;
private final boolean dropDelete;
+ private final String bucketInfo;
public FileRewriteCompactTask(
CompactRewriter rewriter,
CompactUnit unit,
boolean dropDelete,
- @Nullable CompactionMetrics.Reporter metricsReporter) {
+ @Nullable CompactionMetrics.Reporter metricsReporter,
+ String bucketInfo) {
super(metricsReporter);
this.rewriter = rewriter;
this.outputLevel = unit.outputLevel();
this.files = unit.files();
this.dropDelete = dropDelete;
+ this.bucketInfo = bucketInfo;
}
@Override
diff --git
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java
index 6ee834acdd..59e5ca6828 100644
---
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java
+++
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManager.java
@@ -69,11 +69,10 @@ public class MergeTreeCompactManager extends
CompactFutureManager {
private final boolean needLookup;
private final boolean forceRewriteAllFiles;
private final boolean forceKeepDelete;
+ private final String bucketInfo;
@Nullable private final RecordLevelExpire recordLevelExpire;
- private String logInfo = "";
-
public MergeTreeCompactManager(
ExecutorService executor,
Levels levels,
@@ -88,7 +87,8 @@ public class MergeTreeCompactManager extends
CompactFutureManager {
boolean needLookup,
@Nullable RecordLevelExpire recordLevelExpire,
boolean forceRewriteAllFiles,
- boolean forceKeepDelete) {
+ boolean forceKeepDelete,
+ String bucketInfo) {
this.executor = executor;
this.levels = levels;
this.strategy = strategy;
@@ -103,15 +103,11 @@ public class MergeTreeCompactManager extends
CompactFutureManager {
this.needLookup = needLookup;
this.forceRewriteAllFiles = forceRewriteAllFiles;
this.forceKeepDelete = forceKeepDelete;
+ this.bucketInfo = bucketInfo;
MetricUtils.safeCall(this::reportMetrics, LOG);
}
- /** Set additional compact task information for logging purposes. */
- public void setLogInfo(String logInfo) {
- this.logInfo = logInfo;
- }
-
@Override
public boolean shouldWaitForLatestCompaction() {
return levels.numberOfSortedRuns() > numSortedRunStopTrigger;
@@ -223,7 +219,9 @@ public class MergeTreeCompactManager extends
CompactFutureManager {
CompactTask task;
if (unit.fileRewrite()) {
- task = new FileRewriteCompactTask(rewriter, unit, dropDelete,
metricsReporter);
+ task =
+ new FileRewriteCompactTask(
+ rewriter, unit, dropDelete, metricsReporter,
bucketInfo);
} else {
task =
new MergeTreeCompactTask(
@@ -236,7 +234,8 @@ public class MergeTreeCompactManager extends
CompactFutureManager {
metricsReporter,
compactDfSupplier,
recordLevelExpire,
- forceRewriteAllFiles);
+ forceRewriteAllFiles,
+ bucketInfo);
}
if (LOG.isDebugEnabled()) {
@@ -251,7 +250,6 @@ public class MergeTreeCompactManager extends
CompactFutureManager {
file.fileName(),
file.level(), file.fileSize()))
.collect(Collectors.joining(", ")));
}
- task.setLogInfo(logInfo);
taskFuture = executor.submit(task);
if (metricsReporter != null) {
metricsReporter.increaseCompactionsQueuedCount();
diff --git
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
index 23d997c0cf..d90fd76c9f 100644
---
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
+++
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
@@ -173,34 +173,28 @@ public class MergeTreeCompactManagerFactory implements
KvCompactionManagerFactor
if (metricsReporter != null) {
rewriter.setMetricsReporter(metricsReporter);
}
- MergeTreeCompactManager compactManager =
- new MergeTreeCompactManager(
- compactExecutor,
- levels,
- compactStrategy,
- keyComparator,
- options.compactionFileSize(true),
- options.numSortedRunStopTrigger(),
- rewriter,
- metricsReporter,
- dvMaintainer,
- options.prepareCommitWaitCompaction(),
- options.needLookup(),
- recordLevelExpire,
- options.forceRewriteAllFiles(),
- options.isChainTable());
- compactManager.setLogInfo(compactTaskLogInfo(partition, bucket));
- return compactManager;
- }
-
- private String compactTaskLogInfo(BinaryRow partition, int bucket) {
- String partitionString;
- try {
- partitionString =
readerFactoryBuilder.pathFactory().getPartitionString(partition);
- } catch (Exception e) {
- partitionString = partition.toString();
+ String bucketInfo = "bucket=" + bucket;
+ if (partition.getFieldCount() > 0) {
+ String partitionString =
+
readerFactoryBuilder.pathFactory().getPartitionString(partition);
+ bucketInfo = String.format("partition=%s, ", partitionString) +
bucketInfo;
}
- return String.format("partition=%s, bucket=%d", partitionString,
bucket);
+ return new MergeTreeCompactManager(
+ compactExecutor,
+ levels,
+ compactStrategy,
+ keyComparator,
+ options.compactionFileSize(true),
+ options.numSortedRunStopTrigger(),
+ rewriter,
+ metricsReporter,
+ dvMaintainer,
+ options.prepareCommitWaitCompaction(),
+ options.needLookup(),
+ recordLevelExpire,
+ options.forceRewriteAllFiles(),
+ options.isChainTable(),
+ bucketInfo);
}
private CompactStrategy createCompactStrategy(
diff --git
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactTask.java
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactTask.java
index 667a965c62..db6d8e23e8 100644
---
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactTask.java
+++
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactTask.java
@@ -63,8 +63,9 @@ public class MergeTreeCompactTask extends CompactTask {
@Nullable CompactionMetrics.Reporter metricsReporter,
Supplier<CompactDeletionFile> compactDfSupplier,
@Nullable RecordLevelExpire recordLevelExpire,
- boolean forceRewriteAllFiles) {
- super(metricsReporter);
+ boolean forceRewriteAllFiles,
+ String bucketInfo) {
+ super(metricsReporter, bucketInfo);
this.minFileSize = minFileSize;
this.rewriter = rewriter;
this.outputLevel = unit.outputLevel();
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 2961a37029..fb1db4e841 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
@@ -62,6 +62,7 @@ import java.util.OptionalLong;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.function.Function;
+import java.util.function.Supplier;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
@@ -540,7 +541,16 @@ public abstract class AbstractFileStoreWrite<T> implements
FileStoreWrite<T> {
}
private RestoreFiles scanExistingFileMetas(BinaryRow partition, int
bucket) {
- String partInfo = partitionInfo(partition);
+ Supplier<String> partInfo =
+ () ->
+ partitionType.getFieldCount() > 0
+ ? "partition "
+ + getPartitionComputer(
+ partitionType,
+
PARTITION_DEFAULT_NAME.defaultValue(),
+ legacyPartitionName)
+ .generatePartValues(partition)
+ : "table";
RestoreFiles restored;
try {
restored =
@@ -553,7 +563,7 @@ public abstract class AbstractFileStoreWrite<T> implements
FileStoreWrite<T> {
throw new RuntimeException(
String.format(
"Failed to restore existing files for %s, bucket
%d.",
- partInfo, bucket),
+ partInfo.get(), bucket),
e);
}
Integer restoredTotalBuckets = restored.totalBuckets();
@@ -566,22 +576,11 @@ public abstract class AbstractFileStoreWrite<T>
implements FileStoreWrite<T> {
String.format(
"Try to write %s with a new bucket num %d, but the
previous bucket num is %d. "
+ "Please switch to batch mode, and
perform INSERT OVERWRITE to rescale current data layout first.",
- partInfo, numBuckets, totalBuckets));
+ partInfo.get(), numBuckets, totalBuckets));
}
return restored;
}
- private String partitionInfo(BinaryRow partition) {
- return partitionType.getFieldCount() > 0
- ? "partition "
- + getPartitionComputer(
- partitionType,
- PARTITION_DEFAULT_NAME.defaultValue(),
- legacyPartitionName)
- .generatePartValues(partition)
- : "table";
- }
-
private ExecutorService compactExecutor() {
if (lazyCompactExecutor == null) {
lazyCompactExecutor =
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/WriteRestore.java
b/paimon-core/src/main/java/org/apache/paimon/operation/WriteRestore.java
index 0d8db17313..5d4e335571 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/WriteRestore.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/WriteRestore.java
@@ -39,20 +39,13 @@ public interface WriteRestore {
@Nullable
static Integer extractDataFiles(List<ManifestEntry> entries,
List<DataFileMeta> dataFiles) {
- return extractDataFiles(entries, dataFiles, null);
- }
-
- @Nullable
- static Integer extractDataFiles(
- List<ManifestEntry> entries, List<DataFileMeta> dataFiles,
@Nullable String context) {
Integer totalBuckets = null;
for (ManifestEntry entry : entries) {
if (totalBuckets != null && totalBuckets != entry.totalBuckets()) {
- String contextInfo = context == null ? "" : " for " + context;
throw new RuntimeException(
String.format(
- "Bucket data files%s has different total
bucket number, %s vs %s, this should be a bug.",
- contextInfo, totalBuckets,
entry.totalBuckets()));
+ "Bucket data files has different total bucket
number, %s vs %s, this should be a bug.",
+ totalBuckets, entry.totalBuckets()));
}
totalBuckets = entry.totalBuckets();
dataFiles.add(entry.file());
diff --git
a/paimon-core/src/test/java/org/apache/paimon/mergetree/MergeTreeTestBase.java
b/paimon-core/src/test/java/org/apache/paimon/mergetree/MergeTreeTestBase.java
index 9769e59e17..67e787770f 100644
---
a/paimon-core/src/test/java/org/apache/paimon/mergetree/MergeTreeTestBase.java
+++
b/paimon-core/src/test/java/org/apache/paimon/mergetree/MergeTreeTestBase.java
@@ -456,7 +456,8 @@ public abstract class MergeTreeTestBase {
options.needLookup(),
null,
false,
- false);
+ false,
+ "");
}
static class MockFailResultCompactionManager extends
MergeTreeCompactManager {
@@ -482,7 +483,8 @@ public abstract class MergeTreeTestBase {
false,
null,
false,
- false);
+ false,
+ "");
}
protected CompactResult obtainCompactResult()
diff --git
a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerTest.java
b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerTest.java
index 07406baac5..05c041521b 100644
---
a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerTest.java
@@ -243,7 +243,8 @@ public class MergeTreeCompactManagerTest {
true,
null,
false,
- false);
+ false,
+ "");
MergeTreeCompactManager defaultManager =
new MergeTreeCompactManager(
@@ -260,7 +261,8 @@ public class MergeTreeCompactManagerTest {
false,
null,
false,
- false);
+ false,
+ "");
assertThat(lookupManager.compactNotCompleted()).isTrue();
assertThat(defaultManager.compactNotCompleted()).isFalse();
@@ -403,7 +405,8 @@ public class MergeTreeCompactManagerTest {
false,
null,
false,
- true); // keepDelete=true
+ true,
+ ""); // keepDelete=true
try {
manager.triggerCompaction(false);
@@ -527,7 +530,8 @@ public class MergeTreeCompactManagerTest {
false,
null,
false,
- false);
+ false,
+ "");
manager.triggerCompaction(false);
manager.getCompactionResult(true);
List<LevelMinMax> outputs =
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/coordinator/TableWriteCoordinator.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/coordinator/TableWriteCoordinator.java
index 4ca6228124..1cff11eeff 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/coordinator/TableWriteCoordinator.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/coordinator/TableWriteCoordinator.java
@@ -28,7 +28,6 @@ import org.apache.paimon.manifest.ManifestEntry;
import org.apache.paimon.operation.FileStoreScan;
import org.apache.paimon.operation.WriteRestore;
import org.apache.paimon.table.FileStoreTable;
-import org.apache.paimon.utils.FileStorePathFactory;
import
org.apache.paimon.shade.caffeine2.com.github.benmanes.caffeine.cache.Cache;
import
org.apache.paimon.shade.caffeine2.com.github.benmanes.caffeine.cache.Caffeine;
@@ -43,7 +42,6 @@ import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
-import static org.apache.paimon.CoreOptions.PARTITION_DEFAULT_NAME;
import static
org.apache.paimon.deletionvectors.DeletionVectorsIndexFile.DELETION_VECTORS_INDEX;
import static org.apache.paimon.utils.InstantiationUtil.deserializeObject;
import static org.apache.paimon.utils.InstantiationUtil.serializeObject;
@@ -171,11 +169,7 @@ public class TableWriteCoordinator {
List<DataFileMeta> restoreFiles = new ArrayList<>();
List<ManifestEntry> entries = scan.withPartitionBucket(partition,
bucket).plan().files();
- Integer totalBuckets =
- WriteRestore.extractDataFiles(
- entries,
- restoreFiles,
- String.format("%s, bucket %d",
partitionInfo(partition), bucket));
+ Integer totalBuckets = WriteRestore.extractDataFiles(entries,
restoreFiles);
IndexFileMeta dynamicBucketIndex = null;
if (request.scanDynamicBucketIndex()) {
@@ -218,18 +212,6 @@ public class TableWriteCoordinator {
latestCommittedIdentifiers.clear();
}
- private String partitionInfo(BinaryRow partition) {
- if (table.schema().logicalPartitionType().getFieldCount() == 0) {
- return "table";
- }
- return "partition "
- + FileStorePathFactory.getPartitionComputer(
- table.schema().logicalPartitionType(),
-
table.coreOptions().toConfiguration().get(PARTITION_DEFAULT_NAME),
- table.coreOptions().legacyPartitionName())
- .generatePartValues(partition);
- }
-
private static class CoordinationKey {
private final byte[] content;