This is an automated email from the ASF dual-hosted git repository.

zhoujinsong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/amoro.git


The following commit(s) were added to refs/heads/master by this push:
     new 498ea2716 [AMORO-4263] Persist cleanup operation metrics for Dashboard 
display (#4265)
498ea2716 is described below

commit 498ea2716df26ce0677e53401fb8d5c3058aa2c6
Author: WenLingzhang <[email protected]>
AuthorDate: Wed Jul 8 19:41:59 2026 +0800

    [AMORO-4263] Persist cleanup operation metrics for Dashboard display (#4265)
    
    * Persist cleanup operation metrics for Dashboard display
    
    * fixup  metric
    
    ---------
    
    Co-authored-by: 张文领 <[email protected]>
    Co-authored-by: ConradJam <[email protected]>
---
 .../server/process/DefaultTableProcessStore.java   |  50 +++++-----
 .../DanglingDeleteFilesCleaningProcess.java        |   6 +-
 .../process/iceberg/DataExpiringProcess.java       |   6 +-
 .../iceberg/OrphanFilesCleaningProcess.java        |   6 +-
 .../process/iceberg/SnapshotsExpiringProcess.java  |  10 +-
 .../apache/amoro/maintainer/TableMaintainer.java   |  14 ++-
 .../iceberg/maintainer/IcebergTableMaintainer.java | 101 ++++++++++++++-------
 .../iceberg/maintainer/MixedTableMaintainer.java   |  74 ++++++++++++---
 8 files changed, 185 insertions(+), 82 deletions(-)

diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/DefaultTableProcessStore.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/DefaultTableProcessStore.java
index a57d1db21..0dfce34f9 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/DefaultTableProcessStore.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/DefaultTableProcessStore.java
@@ -363,19 +363,29 @@ public class DefaultTableProcessStore extends 
PersistentBase implements TablePro
 
       switch (processEvent) {
         case SUBMIT_REQUESTED:
-          updateExternalProcessIdentifier(newStatus, 
externalProcessIdentifier);
+          if (isTerminal(newStatus)) {
+            // Process completed before we could record SUBMITTED status.
+            // Save summary now — COMPLETE_SUCCESS will be rejected by 
validTransition.
+            begin()
+                .updateTableProcessStatus(newStatus)
+                .updateExternalProcessIdentifier(externalProcessIdentifier)
+                .updateSummary(summary)
+                .commit();
+          } else {
+            updateExternalProcessIdentifier(newStatus, 
externalProcessIdentifier);
+          }
           break;
         case COMPLETE_SUCCESS:
-          updateTableProcessStatus(newStatus);
+          updateTableProcessStatus(newStatus, "", summary);
           break;
         case COMPLETE_FAILED:
-          updateTableProcessStatus(newStatus, reason);
+          updateTableProcessStatus(newStatus, reason, summary);
           break;
         case CANCEL_REQUESTED:
-          updateTableProcessStatus(newStatus, reason);
+          updateTableProcessStatus(newStatus, reason, summary);
           break;
         case KILL_REQUESTED:
-          updateTableProcessStatus(newStatus, reason);
+          updateTableProcessStatus(newStatus, reason, summary);
           break;
         case RETRY_REQUESTED:
           updateTableProcessRetryTimes(getRetryNumber() + 1);
@@ -390,21 +400,14 @@ public class DefaultTableProcessStore extends 
PersistentBase implements TablePro
   }
 
   /**
-   * Update status and persist without message.
-   *
-   * @param status new status
-   */
-  private void updateTableProcessStatus(ProcessStatus status) {
-    updateTableProcessStatus(status, "");
-  }
-
-  /**
-   * Update status with message and persist.
+   * Update status with message and summary, then persist.
    *
    * @param status new status
    * @param message fail or info message
+   * @param summary summary map
    */
-  private void updateTableProcessStatus(ProcessStatus status, String message) {
+  private void updateTableProcessStatus(
+      ProcessStatus status, String message, Map<String, String> summary) {
     switch (status) {
       case SUBMITTED:
       case RUNNING:
@@ -418,6 +421,7 @@ public class DefaultTableProcessStore extends 
PersistentBase implements TablePro
         begin()
             .updateTableProcessStatus(status)
             .updateFinishTime(System.currentTimeMillis())
+            .updateSummary(summary)
             .commit();
         break;
       case FAILED:
@@ -425,6 +429,7 @@ public class DefaultTableProcessStore extends 
PersistentBase implements TablePro
             .updateTableProcessStatus(status)
             .updateTableProcessFailMessage(message)
             .updateFinishTime(System.currentTimeMillis())
+            .updateSummary(summary)
             .commit();
         break;
       default:
@@ -585,11 +590,13 @@ public class DefaultTableProcessStore extends 
PersistentBase implements TablePro
      */
     @Override
     public TableProcessOperation updateSummary(Map<String, String> summary) {
-      operations.add(
-          () -> {
-            oldMeta.setSummary(summary);
-          });
-      metaOperation = true;
+      if (summary != null) {
+        operations.add(
+            () -> {
+              oldMeta.setSummary(summary);
+            });
+        metaOperation = true;
+      }
       return this;
     }
 
@@ -656,6 +663,7 @@ public class DefaultTableProcessStore extends 
PersistentBase implements TablePro
         meta.setRetryNumber(oldMeta.getRetryNumber());
         meta.setCreateTime(oldMeta.getCreateTime());
         meta.setFinishTime(oldMeta.getFinishTime());
+        meta.setSummary(oldMeta.getSummary());
       }
     }
   }
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/DanglingDeleteFilesCleaningProcess.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/DanglingDeleteFilesCleaningProcess.java
index 5c5fef19c..4112ccebc 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/DanglingDeleteFilesCleaningProcess.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/DanglingDeleteFilesCleaningProcess.java
@@ -40,6 +40,8 @@ public class DanglingDeleteFilesCleaningProcess extends 
TableProcess implements
   private static final Logger LOG =
       LoggerFactory.getLogger(DanglingDeleteFilesCleaningProcess.class);
 
+  private volatile Map<String, String> summary = Maps.newLinkedHashMap();
+
   public DanglingDeleteFilesCleaningProcess(TableRuntime tableRuntime, 
ExecuteEngine engine) {
     super(tableRuntime, engine);
   }
@@ -54,7 +56,7 @@ public class DanglingDeleteFilesCleaningProcess extends 
TableProcess implements
     try {
       AmoroTable<?> amoroTable = tableRuntime.loadTable();
       TableMaintainer tableMaintainer = 
TableMaintainerFactory.create(amoroTable, tableRuntime);
-      tableMaintainer.cleanDanglingDeleteFiles();
+      summary = tableMaintainer.cleanDanglingDeleteFiles();
       tableRuntime.updateState(
           DefaultTableRuntime.CLEANUP_STATE_KEY,
           cleanUp -> 
cleanUp.setLastDanglingDeleteFilesCleanTime(System.currentTimeMillis()));
@@ -79,6 +81,6 @@ public class DanglingDeleteFilesCleaningProcess extends 
TableProcess implements
 
   @Override
   public Map<String, String> getSummary() {
-    return Maps.newHashMap();
+    return summary;
   }
 }
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/DataExpiringProcess.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/DataExpiringProcess.java
index d3f9c0cc1..d1dc17fec 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/DataExpiringProcess.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/DataExpiringProcess.java
@@ -39,6 +39,8 @@ public class DataExpiringProcess extends TableProcess 
implements LocalProcess {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(DataExpiringProcess.class);
 
+  private volatile Map<String, String> summary = Maps.newLinkedHashMap();
+
   public DataExpiringProcess(TableRuntime tableRuntime, ExecuteEngine engine) {
     super(tableRuntime, engine);
   }
@@ -53,7 +55,7 @@ public class DataExpiringProcess extends TableProcess 
implements LocalProcess {
     try {
       AmoroTable<?> amoroTable = tableRuntime.loadTable();
       TableMaintainer tableMaintainer = 
TableMaintainerFactory.create(amoroTable, tableRuntime);
-      tableMaintainer.expireData();
+      summary = tableMaintainer.expireData();
       tableRuntime.updateState(
           DefaultTableRuntime.CLEANUP_STATE_KEY,
           cleanUp -> 
cleanUp.setLastDataExpiringTime(System.currentTimeMillis()));
@@ -75,6 +77,6 @@ public class DataExpiringProcess extends TableProcess 
implements LocalProcess {
 
   @Override
   public Map<String, String> getSummary() {
-    return Maps.newHashMap();
+    return summary;
   }
 }
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/OrphanFilesCleaningProcess.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/OrphanFilesCleaningProcess.java
index 3c04c58b1..2185c16ab 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/OrphanFilesCleaningProcess.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/OrphanFilesCleaningProcess.java
@@ -39,6 +39,8 @@ public class OrphanFilesCleaningProcess extends TableProcess 
implements LocalPro
 
   private static final Logger LOG = 
LoggerFactory.getLogger(OrphanFilesCleaningProcess.class);
 
+  private volatile Map<String, String> summary = Maps.newLinkedHashMap();
+
   public OrphanFilesCleaningProcess(TableRuntime tableRuntime, ExecuteEngine 
engine) {
     super(tableRuntime, engine);
   }
@@ -53,7 +55,7 @@ public class OrphanFilesCleaningProcess extends TableProcess 
implements LocalPro
     try {
       AmoroTable<?> amoroTable = tableRuntime.loadTable();
       TableMaintainer tableMaintainer = 
TableMaintainerFactory.create(amoroTable, tableRuntime);
-      tableMaintainer.cleanOrphanFiles();
+      summary = tableMaintainer.cleanOrphanFiles();
       tableRuntime.updateState(
           DefaultTableRuntime.CLEANUP_STATE_KEY,
           cleanUp -> 
cleanUp.setLastOrphanFilesCleanTime(System.currentTimeMillis()));
@@ -75,6 +77,6 @@ public class OrphanFilesCleaningProcess extends TableProcess 
implements LocalPro
 
   @Override
   public Map<String, String> getSummary() {
-    return Maps.newHashMap();
+    return summary;
   }
 }
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/SnapshotsExpiringProcess.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/SnapshotsExpiringProcess.java
index e44e8b2ca..6635209e3 100755
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/SnapshotsExpiringProcess.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/process/iceberg/SnapshotsExpiringProcess.java
@@ -26,7 +26,7 @@ import org.apache.amoro.maintainer.TableMaintainer;
 import org.apache.amoro.process.ExecuteEngine;
 import org.apache.amoro.process.LocalProcess;
 import org.apache.amoro.process.TableProcess;
-import org.apache.amoro.server.optimizing.maintainer.TableMaintainers;
+import org.apache.amoro.server.optimizing.maintainer.TableMaintainerFactory;
 import org.apache.amoro.server.table.DefaultTableRuntime;
 import org.apache.amoro.shade.guava32.com.google.common.collect.Maps;
 import org.slf4j.Logger;
@@ -39,6 +39,8 @@ public class SnapshotsExpiringProcess extends TableProcess 
implements LocalProce
 
   private static final Logger LOG = 
LoggerFactory.getLogger(SnapshotsExpiringProcess.class);
 
+  private volatile Map<String, String> summary = Maps.newLinkedHashMap();
+
   public SnapshotsExpiringProcess(TableRuntime tableRuntime, ExecuteEngine 
engine) {
     super(tableRuntime, engine);
   }
@@ -52,8 +54,8 @@ public class SnapshotsExpiringProcess extends TableProcess 
implements LocalProce
   public void run() {
     try {
       AmoroTable<?> amoroTable = tableRuntime.loadTable();
-      TableMaintainer tableMaintainer = TableMaintainers.create(amoroTable, 
tableRuntime);
-      tableMaintainer.expireSnapshots();
+      TableMaintainer tableMaintainer = 
TableMaintainerFactory.create(amoroTable, tableRuntime);
+      summary = tableMaintainer.expireSnapshots();
       tableRuntime.updateState(
           DefaultTableRuntime.CLEANUP_STATE_KEY,
           cleanUp -> 
cleanUp.setLastSnapshotsExpiringTime(System.currentTimeMillis()));
@@ -75,6 +77,6 @@ public class SnapshotsExpiringProcess extends TableProcess 
implements LocalProce
 
   @Override
   public Map<String, String> getSummary() {
-    return Maps.newHashMap();
+    return summary;
   }
 }
diff --git 
a/amoro-common/src/main/java/org/apache/amoro/maintainer/TableMaintainer.java 
b/amoro-common/src/main/java/org/apache/amoro/maintainer/TableMaintainer.java
index 6719e1507..cb94f5b9c 100644
--- 
a/amoro-common/src/main/java/org/apache/amoro/maintainer/TableMaintainer.java
+++ 
b/amoro-common/src/main/java/org/apache/amoro/maintainer/TableMaintainer.java
@@ -18,6 +18,10 @@
 
 package org.apache.amoro.maintainer;
 
+import org.apache.amoro.shade.guava32.com.google.common.collect.Maps;
+
+import java.util.Map;
+
 /**
  * API for maintaining table.
  *
@@ -27,24 +31,24 @@ package org.apache.amoro.maintainer;
 public interface TableMaintainer {
 
   /** Clean table orphan files. Includes: data files, metadata files. */
-  void cleanOrphanFiles();
+  Map<String, String> cleanOrphanFiles();
 
   /** Clean table dangling delete files. */
-  default void cleanDanglingDeleteFiles() {
-    // DO nothing by default
+  default Map<String, String> cleanDanglingDeleteFiles() {
+    return Maps.newHashMap();
   }
 
   /**
    * Expire snapshots. The optimizing based on the snapshot that the current 
table relies on will
    * not expire according to TableRuntime.
    */
-  void expireSnapshots();
+  Map<String, String> expireSnapshots();
 
   /**
    * Expire historical data based on the expiration field, and data that 
exceeds the retention
    * period will be purged
    */
-  void expireData();
+  Map<String, String> expireData();
 
   /** Auto create tags for table. */
   void autoCreateTags();
diff --git 
a/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/IcebergTableMaintainer.java
 
b/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/IcebergTableMaintainer.java
index f121d1fe6..68d90e3a1 100644
--- 
a/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/IcebergTableMaintainer.java
+++ 
b/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/IcebergTableMaintainer.java
@@ -36,6 +36,7 @@ import org.apache.amoro.maintainer.TableMaintainer;
 import org.apache.amoro.maintainer.TableMaintainerContext;
 import 
org.apache.amoro.shade.guava32.com.google.common.annotations.VisibleForTesting;
 import org.apache.amoro.shade.guava32.com.google.common.base.Strings;
+import org.apache.amoro.shade.guava32.com.google.common.collect.ImmutableMap;
 import org.apache.amoro.shade.guava32.com.google.common.collect.Iterables;
 import org.apache.amoro.shade.guava32.com.google.common.collect.Maps;
 import org.apache.amoro.shade.guava32.com.google.common.collect.Sets;
@@ -134,56 +135,68 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
   }
 
   @Override
-  public void cleanOrphanFiles() {
+  public Map<String, String> cleanOrphanFiles() {
     TableConfiguration tableConfiguration = context.getTableConfiguration();
     MaintainerMetrics metrics = context.getMetrics();
 
     if (!tableConfiguration.isCleanOrphanEnabled()) {
-      return;
+      return Maps.newHashMap();
     }
 
     long keepTime = tableConfiguration.getOrphanExistingMinutes() * 60 * 1000;
 
-    cleanContentFiles(System.currentTimeMillis() - keepTime, metrics);
+    int dataDeleted = cleanContentFiles(System.currentTimeMillis() - keepTime, 
metrics);
 
     // refresh
     table.refresh();
 
     // clear metadata files
-    cleanMetadata(System.currentTimeMillis() - keepTime, metrics);
+    int metadataDeleted = cleanMetadata(System.currentTimeMillis() - keepTime, 
metrics);
+
+    Map<String, String> summary = Maps.newLinkedHashMap();
+    summary.put("orphan-data-files-cleaned", String.valueOf(dataDeleted));
+    summary.put("orphan-metadata-files-cleaned", 
String.valueOf(metadataDeleted));
+    return summary;
   }
 
   @Override
-  public void cleanDanglingDeleteFiles() {
+  public Map<String, String> cleanDanglingDeleteFiles() {
     TableConfiguration tableConfiguration = context.getTableConfiguration();
     if (!tableConfiguration.isDeleteDanglingDeleteFilesEnabled()) {
-      return;
+      return Maps.newHashMap();
     }
 
     Snapshot currentSnapshot = table.currentSnapshot();
     if (currentSnapshot == null) {
-      return;
+      return ImmutableMap.of(
+          "dangling-delete-files-cleaned", "0", 
"dangling-delete-files-size-bytes", "0");
     }
     Optional<String> totalDeleteFiles =
         
Optional.ofNullable(currentSnapshot.summary().get(SnapshotSummary.TOTAL_DELETE_FILES_PROP));
     if (totalDeleteFiles.isPresent() && Long.parseLong(totalDeleteFiles.get()) 
> 0) {
       // clear dangling delete files
-      doCleanDanglingDeleteFiles();
+      return doCleanDanglingDeleteFiles();
     } else {
       LOG.debug(
           "There are no delete files here, so there is no need to clean 
dangling delete file for table {}",
           table.name());
+      return ImmutableMap.of(
+          "dangling-delete-files-cleaned", "0", 
"dangling-delete-files-size-bytes", "0");
     }
   }
 
   @Override
-  public void expireSnapshots() {
+  public Map<String, String> expireSnapshots() {
     if (!expireSnapshotEnabled()) {
-      return;
+      return Maps.newHashMap();
     }
-    expireSnapshots(
-        mustOlderThan(System.currentTimeMillis()),
-        context.getTableConfiguration().getSnapshotMinCount());
+    int cleaned =
+        expireSnapshots(
+            mustOlderThan(System.currentTimeMillis()),
+            context.getTableConfiguration().getSnapshotMinCount(),
+            expireSnapshotNeedToExcludeFiles());
+
+    return ImmutableMap.of("snapshot-files-cleaned", String.valueOf(cleaned));
   }
 
   public boolean expireSnapshotEnabled() {
@@ -196,7 +209,7 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
     expireSnapshots(mustOlderThan, minCount, 
expireSnapshotNeedToExcludeFiles());
   }
 
-  private void expireSnapshots(long olderThan, int minCount, Set<String> 
exclude) {
+  protected int expireSnapshots(long olderThan, int minCount, Set<String> 
exclude) {
     LOG.debug(
         "Starting snapshots expiration for table {}, expiring snapshots older 
than {} and retain last {} snapshots, excluding {}",
         table.name(),
@@ -225,11 +238,12 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
 
     int collectedFiles = expiredFileCleaner.fileCount();
     expiredFileCleaner.clear();
+    int cleanedFiles = expiredFileCleaner.cleanedFileCount();
     if (collectedFiles > 0) {
       LOG.info(
           "Expired {}/{} files for table {} order than {}",
           collectedFiles,
-          expiredFileCleaner.cleanedFileCount(),
+          cleanedFiles,
           table.name(),
           DateTimeUtil.formatTimestampMillis(olderThan));
     } else {
@@ -238,6 +252,7 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
           table.name(),
           DateTimeUtil.formatTimestampMillis(olderThan));
     }
+    return cleanedFiles;
   }
 
   /**
@@ -263,17 +278,18 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
   }
 
   @Override
-  public void expireData() {
+  public Map<String, String> expireData() {
     DataExpirationConfig expirationConfig = 
context.getTableConfiguration().getExpiringDataConfig();
     try {
       Types.NestedField field = 
table.schema().findField(expirationConfig.getExpirationField());
       if (!isValidDataExpirationField(expirationConfig, field, table.name())) {
-        return;
+        return Maps.newHashMap();
       }
 
-      expireDataFrom(expirationConfig, expireBaseOnRule(expirationConfig, 
field));
+      return expireDataFrom(expirationConfig, 
expireBaseOnRule(expirationConfig, field));
     } catch (Throwable t) {
       LOG.error("Unexpected purge error for table {} ", tableIdentifier, t);
+      return Maps.newHashMap();
     }
   }
 
@@ -304,9 +320,10 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
    *     zone
    */
   @VisibleForTesting
-  public void expireDataFrom(DataExpirationConfig expirationConfig, Instant 
instant) {
+  public Map<String, String> expireDataFrom(
+      DataExpirationConfig expirationConfig, Instant instant) {
     if (instant.equals(Instant.MIN)) {
-      return;
+      return Maps.newHashMap();
     }
 
     long expireTimestamp = 
instant.minusMillis(expirationConfig.getRetentionTime()).toEpochMilli();
@@ -322,6 +339,18 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
 
     ExpireFiles expiredFiles = expiredFileScan(expirationConfig, dataFilter, 
expireTimestamp);
     expireFiles(expiredFiles, expireTimestamp);
+
+    int dataCount = expiredFiles.dataFiles.size();
+    int deleteCount = expiredFiles.deleteFiles.size();
+    long dataSize = 
expiredFiles.dataFiles.stream().mapToLong(ContentFile::fileSizeInBytes).sum();
+    long deleteSize =
+        
expiredFiles.deleteFiles.stream().mapToLong(ContentFile::fileSizeInBytes).sum();
+    Map<String, String> summary = Maps.newLinkedHashMap();
+    summary.put("expired-data-files-cleaned", String.valueOf(dataCount));
+    summary.put("expired-data-files-size-bytes", String.valueOf(dataSize));
+    summary.put("expired-delete-files-cleaned", String.valueOf(deleteCount));
+    summary.put("expired-delete-files-size-bytes", String.valueOf(deleteSize));
+    return summary;
   }
 
   @Override
@@ -330,7 +359,7 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
     new AutoCreateIcebergTagAction(table, tagConfiguration, 
LocalDateTime.now()).execute();
   }
 
-  public void cleanContentFiles(long lastTime, MaintainerMetrics metrics) {
+  public int cleanContentFiles(long lastTime, MaintainerMetrics metrics) {
     // For clean data files, should getRuntime valid files in the base store 
and the change store,
     // so acquire in advance
     // to prevent repeated acquisition
@@ -339,20 +368,21 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
         "Starting cleaning orphan content files for table {} before {}",
         table.name(),
         Instant.ofEpochMilli(lastTime));
-    clearInternalTableContentsFiles(lastTime, validFiles, metrics);
+    return clearInternalTableContentsFiles(lastTime, validFiles, metrics);
   }
 
-  public void cleanMetadata(long lastTime, MaintainerMetrics metrics) {
+  public int cleanMetadata(long lastTime, MaintainerMetrics metrics) {
     LOG.info(
         "Starting cleaning metadata files for table {} before {}",
         table.name(),
         Instant.ofEpochMilli(lastTime));
-    clearInternalTableMetadata(lastTime, metrics);
+    return clearInternalTableMetadata(lastTime, metrics);
   }
 
-  public void doCleanDanglingDeleteFiles() {
+  public Map<String, String> doCleanDanglingDeleteFiles() {
     LOG.info("Starting cleaning dangling delete files for table {}", 
table.name());
-    int danglingDeleteFilesCnt = clearInternalTableDanglingDeleteFiles();
+    Set<DeleteFile> danglingDeleteFiles = 
IcebergTableUtil.getDanglingDeleteFiles(table);
+    int danglingDeleteFilesCnt = 
clearInternalTableDanglingDeleteFiles(danglingDeleteFiles);
     runWithCondition(
         danglingDeleteFilesCnt > 0,
         () ->
@@ -360,6 +390,12 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
                 "Deleted {} dangling delete files for table {}",
                 danglingDeleteFilesCnt,
                 table.name()));
+
+    Map<String, String> summary = Maps.newLinkedHashMap();
+    long totalSize = 
danglingDeleteFiles.stream().mapToLong(ContentFile::fileSizeInBytes).sum();
+    summary.put("dangling-delete-files-cleaned", 
String.valueOf(danglingDeleteFilesCnt));
+    summary.put("dangling-delete-files-size-bytes", String.valueOf(totalSize));
+    return summary;
   }
 
   public long mustOlderThan(long now) {
@@ -402,7 +438,7 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
     return (AuthenticatedFileIO) table.io();
   }
 
-  private void clearInternalTableContentsFiles(
+  private int clearInternalTableContentsFiles(
       long lastTime, Set<String> exclude, MaintainerMetrics metrics) {
     String dataLocation = table.location() + File.separator + DATA_FOLDER_NAME;
     int expected = 0, deleted = 0;
@@ -443,9 +479,10 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
               table.name());
           metrics.recordOrphanDataFilesCleaned(finalExpected, finalDeleted);
         });
+    return finalDeleted;
   }
 
-  private void clearInternalTableMetadata(long lastTime, MaintainerMetrics 
metrics) {
+  private int clearInternalTableMetadata(long lastTime, MaintainerMetrics 
metrics) {
     Set<String> validFiles = getValidMetadataFiles(table);
     LOG.info("Found {} valid metadata files for table {}", validFiles.size(), 
table.name());
     Pattern excludeFileNameRegex = getExcludeFileNameRegex(table);
@@ -457,12 +494,13 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
     LOG.info("start orphan files clean in {}", metadataLocation);
 
     AuthenticatedFileIO io = fileIO();
+    int deleted;
     if (io.supportPrefixOperations()) {
       SupportsPrefixOperations pio = io.asPrefixFileIO();
       Set<String> filesToDelete =
           deleteInvalidMetadataFile(
               pio, metadataLocation, lastTime, validFiles, 
excludeFileNameRegex);
-      int deleted = TableFileUtil.deleteFiles(io, filesToDelete);
+      deleted = TableFileUtil.deleteFiles(io, filesToDelete);
 
       runWithCondition(
           !filesToDelete.isEmpty(),
@@ -475,15 +513,16 @@ public class IcebergTableMaintainer implements 
TableMaintainer {
             metrics.recordOrphanMetadataFilesCleaned(filesToDelete.size(), 
deleted);
           });
     } else {
+      deleted = 0;
       LOG.warn(
           String.format(
               "Table %s doesn't support a fileIo with listDirectory or 
listPrefix, so skip clear files.",
               table.name()));
     }
+    return deleted;
   }
 
-  private int clearInternalTableDanglingDeleteFiles() {
-    Set<DeleteFile> danglingDeleteFiles = 
IcebergTableUtil.getDanglingDeleteFiles(table);
+  private int clearInternalTableDanglingDeleteFiles(Set<DeleteFile> 
danglingDeleteFiles) {
     if (danglingDeleteFiles.isEmpty()) {
       return 0;
     }
diff --git 
a/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/MixedTableMaintainer.java
 
b/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/MixedTableMaintainer.java
index 72e4a6429..159f3d7a4 100644
--- 
a/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/MixedTableMaintainer.java
+++ 
b/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/MixedTableMaintainer.java
@@ -30,6 +30,7 @@ import org.apache.amoro.maintainer.TableMaintainerContext;
 import org.apache.amoro.scan.TableEntriesScan;
 import 
org.apache.amoro.shade.guava32.com.google.common.annotations.VisibleForTesting;
 import org.apache.amoro.shade.guava32.com.google.common.base.Strings;
+import org.apache.amoro.shade.guava32.com.google.common.collect.ImmutableMap;
 import org.apache.amoro.shade.guava32.com.google.common.collect.Iterables;
 import org.apache.amoro.shade.guava32.com.google.common.collect.Maps;
 import org.apache.amoro.shade.guava32.com.google.common.collect.Sets;
@@ -70,6 +71,7 @@ import java.util.Queue;
 import java.util.Set;
 import java.util.concurrent.LinkedTransferQueue;
 import java.util.stream.Collectors;
+import java.util.stream.Stream;
 
 /** Table maintainer for mixed-iceberg and mixed-hive tables. */
 public class MixedTableMaintainer implements TableMaintainer {
@@ -99,24 +101,39 @@ public class MixedTableMaintainer implements 
TableMaintainer {
   }
 
   @Override
-  public void cleanOrphanFiles() {
+  public Map<String, String> cleanOrphanFiles() {
+    Map<String, String> result = Maps.newHashMap();
     if (changeMaintainer != null) {
-      changeMaintainer.cleanOrphanFiles();
+      result.putAll(changeMaintainer.cleanOrphanFiles());
     }
-    baseMaintainer.cleanOrphanFiles();
+    baseMaintainer
+        .cleanOrphanFiles()
+        .forEach(
+            (k, v) ->
+                result.merge(
+                    k, v, (a, b) -> String.valueOf(Long.parseLong(a) + 
Long.parseLong(b))));
+    return result;
   }
 
   @Override
-  public void cleanDanglingDeleteFiles() {
+  public Map<String, String> cleanDanglingDeleteFiles() {
     // Mixed table doesn't support clean dangling delete files
+    return Maps.newHashMap();
   }
 
   @Override
-  public void expireSnapshots() {
+  public Map<String, String> expireSnapshots() {
+    Map<String, String> result = Maps.newHashMap();
     if (changeMaintainer != null) {
-      changeMaintainer.expireSnapshots();
+      result.putAll(changeMaintainer.expireSnapshots());
     }
-    baseMaintainer.expireSnapshots();
+    baseMaintainer
+        .expireSnapshots()
+        .forEach(
+            (k, v) ->
+                result.merge(
+                    k, v, (a, b) -> String.valueOf(Long.parseLong(a) + 
Long.parseLong(b))));
+    return result;
   }
 
   @VisibleForTesting
@@ -128,18 +145,19 @@ public class MixedTableMaintainer implements 
TableMaintainer {
   }
 
   @Override
-  public void expireData() {
+  public Map<String, String> expireData() {
     DataExpirationConfig expirationConfig = 
context.getTableConfiguration().getExpiringDataConfig();
     try {
       Types.NestedField field =
           mixedTable.schema().findField(expirationConfig.getExpirationField());
       if (!isValidDataExpirationField(expirationConfig, field, 
mixedTable.name())) {
-        return;
+        return Maps.newHashMap();
       }
 
-      expireDataFrom(expirationConfig, expireMixedBaseOnRule(expirationConfig, 
field));
+      return expireDataFrom(expirationConfig, 
expireMixedBaseOnRule(expirationConfig, field));
     } catch (Throwable t) {
       LOG.error("Unexpected purge error for table {} ", mixedTable.id(), t);
+      return Maps.newHashMap();
     }
   }
 
@@ -158,9 +176,10 @@ public class MixedTableMaintainer implements 
TableMaintainer {
   }
 
   @VisibleForTesting
-  public void expireDataFrom(DataExpirationConfig expirationConfig, Instant 
instant) {
+  public Map<String, String> expireDataFrom(
+      DataExpirationConfig expirationConfig, Instant instant) {
     if (instant.equals(Instant.MIN)) {
-      return;
+      return Maps.newHashMap();
     }
 
     long expireTimestamp = 
instant.minusMillis(expirationConfig.getRetentionTime()).toEpochMilli();
@@ -180,6 +199,26 @@ public class MixedTableMaintainer implements 
TableMaintainer {
         mixedExpiredFileScan(expirationConfig, dataFilter, expireTimestamp);
 
     expireMixedFiles(mixedExpiredFiles.getLeft(), 
mixedExpiredFiles.getRight(), expireTimestamp);
+
+    IcebergTableMaintainer.ExpireFiles changeFiles = 
mixedExpiredFiles.getLeft();
+    IcebergTableMaintainer.ExpireFiles baseFiles = 
mixedExpiredFiles.getRight();
+    int dataCount = changeFiles.dataFiles.size() + baseFiles.dataFiles.size();
+    int deleteCount = changeFiles.deleteFiles.size() + 
baseFiles.deleteFiles.size();
+    long dataSize =
+        Stream.concat(changeFiles.dataFiles.stream(), 
baseFiles.dataFiles.stream())
+            .mapToLong(ContentFile::fileSizeInBytes)
+            .sum();
+    long deleteSize =
+        Stream.concat(changeFiles.deleteFiles.stream(), 
baseFiles.deleteFiles.stream())
+            .mapToLong(ContentFile::fileSizeInBytes)
+            .sum();
+
+    Map<String, String> summary = Maps.newLinkedHashMap();
+    summary.put("expired-data-files-cleaned", String.valueOf(dataCount));
+    summary.put("expired-data-files-size-bytes", String.valueOf(dataSize));
+    summary.put("expired-delete-files-cleaned", String.valueOf(deleteCount));
+    summary.put("expired-delete-files-size-bytes", String.valueOf(deleteSize));
+    return summary;
   }
 
   private Pair<IcebergTableMaintainer.ExpireFiles, 
IcebergTableMaintainer.ExpireFiles>
@@ -347,13 +386,18 @@ public class MixedTableMaintainer implements 
TableMaintainer {
     }
 
     @Override
-    public void expireSnapshots() {
+    public Map<String, String> expireSnapshots() {
       if (!expireSnapshotEnabled()) {
-        return;
+        return Maps.newHashMap();
       }
       long now = System.currentTimeMillis();
       expireFiles(now - getSnapshotsKeepTime());
-      expireSnapshots(getMustOlderThan(now), 
context.getTableConfiguration().getSnapshotMinCount());
+      int cleaned =
+          expireSnapshots(
+              getMustOlderThan(now),
+              context.getTableConfiguration().getSnapshotMinCount(),
+              expireSnapshotNeedToExcludeFiles());
+      return ImmutableMap.of("snapshot-files-cleaned", 
String.valueOf(cleaned));
     }
 
     private long getSnapshotsKeepTime() {

Reply via email to