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());
}
}