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() {