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;

Reply via email to