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 ea85500237 [spark] Improve compaction logs in CompactProcedure (#9169)
ea85500237 is described below

commit ea8550023735ba9c04112a0c753e416aee58b89b
Author: Zouxxyy <[email protected]>
AuthorDate: Wed Aug 12 11:33:09 2026 +0800

    [spark] Improve compaction logs in CompactProcedure (#9169)
---
 .../paimon/spark/procedure/CompactProcedure.java   | 93 ++++++++++++++++++++--
 1 file changed, 85 insertions(+), 8 deletions(-)

diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
index 8fd213da1b..b3bb7d9307 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
@@ -265,6 +265,19 @@ public class CompactProcedure extends BaseProcedure {
             compactStrategy = clusterIncrementalEnabled ? MINOR : FULL;
         }
         boolean fullCompact = compactStrategy.equalsIgnoreCase(FULL);
+
+        long startMillis = System.currentTimeMillis();
+        LOG.info(
+                "Starting compact on table {}, bucket mode {}, compact 
strategy {}, order type {}, "
+                        + "order by {}, partition filtered {}, partition idle 
time {}.",
+                table.fullName(),
+                bucketMode,
+                compactStrategy,
+                orderType,
+                sortColumns.isEmpty() ? "none" : sortColumns,
+                partitionPredicate != null,
+                partitionIdleTime == null ? "none" : partitionIdleTime);
+
         if (orderType.equals(OrderType.NONE)) {
             JavaSparkContext javaSparkContext = new 
JavaSparkContext(spark().sparkContext());
             switch (bucketMode) {
@@ -311,6 +324,11 @@ public class CompactProcedure extends BaseProcedure {
                                     + " only support unaware-bucket 
append-only table yet.");
             }
         }
+
+        LOG.info(
+                "Finished compact on table {}, cost {} ms.",
+                table.fullName(),
+                System.currentTimeMillis() - startMillis);
         return true;
     }
 
@@ -345,11 +363,18 @@ public class CompactProcedure extends BaseProcedure {
                         .collect(Collectors.toList());
 
         if (partitionBuckets.isEmpty()) {
-            LOG.info("Partition bucket is empty, no compact job to execute.");
+            LOG.info(
+                    "No partition bucket to compact for table {}, skip this 
compact job.",
+                    table.fullName());
             return;
         }
 
         int readParallelism = readParallelism(partitionBuckets, spark());
+        LOG.info(
+                "Starting to compact {} partition buckets of table {} with 
read parallelism {}.",
+                partitionBuckets.size(),
+                table.fullName(),
+                readParallelism);
         BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder();
         JavaRDD<byte[]> commitMessageJavaRDD =
                 javaSparkContext
@@ -450,7 +475,9 @@ public class CompactProcedure extends BaseProcedure {
                             .collect(Collectors.toList());
         }
         if (compactionTasks.isEmpty()) {
-            LOG.info("Task plan is empty, no compact job to execute.");
+            LOG.info(
+                    "No append compact task to execute for table {}, skip this 
compact job.",
+                    table.fullName());
             return;
         }
 
@@ -465,6 +492,11 @@ public class CompactProcedure extends BaseProcedure {
         }
 
         int readParallelism = readParallelism(serializedTasks, spark());
+        LOG.info(
+                "Starting to execute {} append compact tasks of table {} with 
read parallelism {}.",
+                serializedTasks.size(),
+                table.fullName(),
+                readParallelism);
         String commitUser = 
createCommitUser(table.coreOptions().toConfiguration());
         JavaRDD<byte[]> commitMessageJavaRDD =
                 javaSparkContext
@@ -522,6 +554,7 @@ public class CompactProcedure extends BaseProcedure {
         List<DataEvolutionCompactTask> compactionTasks;
         Snapshot snapshot = table.snapshotManager().latestSnapshot();
         if (snapshot == null) {
+            LOG.info("Table {} has no snapshot yet, skip this compact job.", 
table.fullName());
             return;
         }
         DataEvolutionCompactCoordinator compactCoordinator =
@@ -533,9 +566,11 @@ public class CompactProcedure extends BaseProcedure {
                         snapshot);
         CommitMessageSerializer messageSerializerser = new 
CommitMessageSerializer();
         String commitUser = 
createCommitUser(table.coreOptions().toConfiguration());
+        int round = 0;
         try {
             while (true) {
                 compactionTasks = compactCoordinator.plan();
+                round++;
                 if (partitionIdleTime != null) {
                     SnapshotReader snapshotReader = table.newSnapshotReader();
                     if (partitionPredicate != null) {
@@ -562,7 +597,11 @@ public class CompactProcedure extends BaseProcedure {
                                     .collect(Collectors.toList());
                 }
                 if (compactionTasks.isEmpty()) {
-                    LOG.info("Task plan is empty, no compact job to execute.");
+                    LOG.info(
+                            "No compact task planned in round {} for table {}, 
"
+                                    + "continue to scan the next batch.",
+                            round,
+                            table.fullName());
                     continue;
                 }
                 boolean containsMaterializeDeletion = 
containsMaterializeDeletion(compactionTasks);
@@ -579,6 +618,14 @@ public class CompactProcedure extends BaseProcedure {
                 }
 
                 int readParallelism = readParallelism(serializedTasks, 
spark());
+                LOG.info(
+                        "Starting to execute {} data evolution compact tasks 
of table {} in round {} "
+                                + "with read parallelism {}, contains 
materialize deletion {}.",
+                        serializedTasks.size(),
+                        table.fullName(),
+                        round,
+                        readParallelism,
+                        containsMaterializeDeletion);
                 JavaRDD<byte[]> commitMessageJavaRDD =
                         javaSparkContext
                                 .parallelize(serializedTasks, readParallelism)
@@ -621,7 +668,11 @@ public class CompactProcedure extends BaseProcedure {
                 }
             }
         } catch (EndOfScanException e) {
-            LOG.info("Catching EndOfScanException, the compact job is 
finishing.");
+            LOG.info(
+                    "Catching EndOfScanException, the compact job of table {} 
is finishing "
+                            + "after {} plan rounds.",
+                    table.fullName(),
+                    round);
         }
     }
 
@@ -679,7 +730,21 @@ public class CompactProcedure extends BaseProcedure {
             snapshotReader.withPartitionFilter(partitionPredicate);
         }
         Map<BinaryRow, DataSplit[]> packedSplits = 
packForSort(snapshotReader.read().dataSplits());
+        // Build the sorter before the emptiness check on purpose: its 
constructor validates the
+        // order columns, and that validation must keep failing fast even for 
an empty table.
         TableSorter sorter = TableSorter.getSorter(table, orderType, 
sortColumns);
+        if (packedSplits.isEmpty()) {
+            LOG.info(
+                    "No data split to sort compact for table {}, skip this 
compact job.",
+                    table.fullName());
+            return;
+        }
+        LOG.info(
+                "Starting to sort compact {} partitions of table {}, order 
type {}, order by {}.",
+                packedSplits.size(),
+                table.fullName(),
+                orderType,
+                sortColumns);
         Dataset<Row> datasetForWrite =
                 packedSplits.values().stream()
                         .map(
@@ -710,6 +775,11 @@ public class CompactProcedure extends BaseProcedure {
                 new IncrementalClusterManager(table, partitionPredicate);
         Map<BinaryRow, CompactUnit> compactUnits =
                 incrementalClusterManager.createCompactUnits(fullCompaction);
+        LOG.info(
+                "Planned {} compact units to incrementally cluster for table 
{}, full compaction {}.",
+                compactUnits.size(),
+                table.fullName(),
+                fullCompaction);
 
         Map<BinaryRow, Pair<List<DataSplit>, CommitMessage>> partitionSplits =
                 
incrementalClusterManager.toSplitsAndRewriteDvFiles(compactUnits);
@@ -720,13 +790,16 @@ public class CompactProcedure extends BaseProcedure {
                         table,
                         incrementalClusterManager.clusterCurve(),
                         incrementalClusterManager.clusterKeys());
-        LOG.info(
-                "Start to sort in partition, cluster curve is {}, cluster keys 
is {}",
-                incrementalClusterManager.clusterCurve(),
-                incrementalClusterManager.clusterKeys());
 
         CoreOptions.ClusteringIncrementalMode mode =
                 incrementalClusterManager.clusteringIncrementalMode();
+        LOG.info(
+                "Start to sort in partition for table {}, cluster curve is {}, 
cluster keys is {}, "
+                        + "incremental mode is {}",
+                table.fullName(),
+                incrementalClusterManager.clusterCurve(),
+                incrementalClusterManager.clusterKeys(),
+                mode);
 
         Dataset<Row> datasetForWrite =
                 partitionSplits.values().stream()
@@ -811,6 +884,10 @@ public class CompactProcedure extends BaseProcedure {
             }
 
             
writer.commit(JavaConverters.asScalaBuffer(clusterMessages).toSeq());
+        } else {
+            LOG.info(
+                    "No data split to incrementally cluster for table {}, skip 
this compact job.",
+                    table.fullName());
         }
     }
 

Reply via email to