This is an automated email from the ASF dual-hosted git repository. xxubai pushed a commit to branch 0.9.x in repository https://gitbox.apache.org/repos/asf/amoro.git
commit 5f618fb629aca220665a1a2eac85326889615def Author: Wang Tao <[email protected]> AuthorDate: Fri Jul 24 15:45:39 2026 +0800 [AMORO-4277] imporove expire_snapshots performance (#4278) * optimize expire snapshots fix: logical error remove some comments * remove duplicate class and move test case to amoro-format-iceberg * Update amoro-format-iceberg/src/main/java/org/apache/amoro/utils/TableFileUtil.java Co-authored-by: Xu Bai <[email protected]> * fix review comment --------- Co-authored-by: Xu Bai <[email protected]> (cherry picked from commit 7be0710b3e134f889ce1f970604428663c0ca97b) --- .../amoro/server/utils/RollingFileCleaner.java | 138 --------------------- .../formats/iceberg/utils/RollingFileCleaner.java | 41 +++--- .../java/org/apache/amoro/utils/TableFileUtil.java | 18 ++- .../iceberg/utils}/TestRollingFileCleaner.java | 3 +- 4 files changed, 44 insertions(+), 156 deletions(-) diff --git a/amoro-ams/src/main/java/org/apache/amoro/server/utils/RollingFileCleaner.java b/amoro-ams/src/main/java/org/apache/amoro/server/utils/RollingFileCleaner.java deleted file mode 100644 index 821bf6661..000000000 --- a/amoro-ams/src/main/java/org/apache/amoro/server/utils/RollingFileCleaner.java +++ /dev/null @@ -1,138 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.amoro.server.utils; - -import org.apache.amoro.io.AuthenticatedFileIO; -import org.apache.amoro.shade.guava32.com.google.common.collect.Sets; -import org.apache.amoro.utils.TableFileUtil; -import org.apache.hadoop.fs.Path; -import org.apache.iceberg.io.BulkDeletionFailureException; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.net.URI; -import java.util.Set; -import java.util.concurrent.atomic.AtomicInteger; - -public class RollingFileCleaner { - private final Set<String> collectedFiles = Sets.newConcurrentHashSet(); - private final Set<String> excludeFiles; - private final Set<String> parentDirectories = Sets.newConcurrentHashSet(); - - private static final int CLEANED_FILE_GROUP_SIZE = 1_000; - - private final AtomicInteger fileCounter = new AtomicInteger(0); - private final AtomicInteger cleanedFileCounter = new AtomicInteger(0); - - private final AuthenticatedFileIO fileIO; - - private static final Logger LOG = LoggerFactory.getLogger(RollingFileCleaner.class); - - public RollingFileCleaner(AuthenticatedFileIO fileIO, Set<String> excludeFiles) { - this.fileIO = fileIO; - this.excludeFiles = - excludeFiles != null - ? Sets.newConcurrentHashSet(excludeFiles) - : Sets.newConcurrentHashSet(); - } - - public void addFile(String filePath) { - if (isFileExcluded(filePath)) { - LOG.debug("File {} is excluded from cleaning", filePath); - return; - } - - collectedFiles.add(filePath); - String parentDir = TableFileUtil.getParent(filePath); - parentDirectories.add(parentDir); - int currentCount = fileCounter.incrementAndGet(); - - if (currentCount % CLEANED_FILE_GROUP_SIZE == 0) { - doCleanFiles(); - } - } - - private boolean isFileExcluded(String filePath) { - if (excludeFiles.isEmpty()) { - return false; - } - - if (excludeFiles.contains(filePath)) { - return true; - } - - String uriPath = URI.create(filePath).getPath(); - if (excludeFiles.contains(uriPath)) { - return true; - } - - String parentPath = new Path(uriPath).getParent().toString(); - return excludeFiles.contains(parentPath); - } - - private void doCleanFiles() { - if (collectedFiles.isEmpty()) { - return; - } - - if (fileIO.supportBulkOperations()) { - try { - fileIO.asBulkFileIO().deleteFiles(collectedFiles); - cleanedFileCounter.addAndGet(collectedFiles.size()); - } catch (BulkDeletionFailureException e) { - LOG.warn("Failed to delete {} expired files in bulk", e.numberFailedObjects()); - } - } else { - for (String filePath : collectedFiles) { - try { - fileIO.deleteFile(filePath); - cleanedFileCounter.incrementAndGet(); - } catch (Exception e) { - LOG.warn("Failed to delete expired file: {}", filePath, e); - } - } - } - // Try to delete empty parent directories - for (String parentDir : parentDirectories) { - TableFileUtil.deleteEmptyDirectory(fileIO, parentDir, excludeFiles); - } - parentDirectories.clear(); - - LOG.debug("Cleaned expired a file group, total files: {}", collectedFiles.size()); - - collectedFiles.clear(); - } - - public int fileCount() { - return fileCounter.get(); - } - - public int cleanedFileCount() { - return cleanedFileCounter.get(); - } - - public void clear() { - if (!collectedFiles.isEmpty()) { - doCleanFiles(); - } - - collectedFiles.clear(); - excludeFiles.clear(); - } -} diff --git a/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/utils/RollingFileCleaner.java b/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/utils/RollingFileCleaner.java index 8d2e8ac06..3fcdcb5b2 100644 --- a/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/utils/RollingFileCleaner.java +++ b/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/utils/RollingFileCleaner.java @@ -27,8 +27,11 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.net.URI; +import java.util.Comparator; +import java.util.List; import java.util.Set; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; /** Rolling file cleaner for Iceberg table maintenance operations. */ public class RollingFileCleaner { @@ -110,27 +113,34 @@ public class RollingFileCleaner { } } } - // Try to delete empty parent directories. Skipped automatically by - // TableFileUtil.deleteEmptyDirectory when the underlying FileIO is an object store. - if (fileIO.supportFileSystemOperations()) { - for (String parentDir : parentDirectories) { - try { - TableFileUtil.deleteEmptyDirectory(fileIO, parentDir, excludeFiles); - } catch (Exception e) { - LOG.warn("Failed to delete empty parent directory: {}", parentDir, e); - } - } - } LOG.debug("Cleaned expired a file group, total files: {}", collectedFiles.size()); } finally { - // Always drop both buffers so a failure in directory cleanup does not cause the - // already-deleted files to be re-submitted on the next batch. - parentDirectories.clear(); collectedFiles.clear(); } } + private void doCleanParentDirectory() { + // Try to delete empty parent directories. Skipped automatically by + // TableFileUtil.deleteEmptyDirectory when the underlying FileIO is an object store. + if (fileIO.supportFileSystemOperations()) { + List<String> parentDirectoriesSorted = + parentDirectories.stream() + .sorted(Comparator.comparingInt(String::length).reversed()) + .collect(Collectors.toList()); + for (String parentDir : parentDirectoriesSorted) { + try { + TableFileUtil.deleteEmptyDirectory(fileIO, parentDir, excludeFiles, parentDirectories); + } catch (Exception e) { + LOG.warn("Failed to delete empty parent directory: {}", parentDir, e); + } + } + } + // Always drop both buffers so a failure in directory cleanup does not cause the + // already-deleted files to be re-submitted on the next batch. + parentDirectories.clear(); + } + public int fileCount() { return fileCounter.get(); } @@ -143,6 +153,9 @@ public class RollingFileCleaner { if (!collectedFiles.isEmpty()) { doCleanFiles(); } + if (!parentDirectories.isEmpty()) { + doCleanParentDirectory(); + } collectedFiles.clear(); excludeFiles.clear(); diff --git a/amoro-format-iceberg/src/main/java/org/apache/amoro/utils/TableFileUtil.java b/amoro-format-iceberg/src/main/java/org/apache/amoro/utils/TableFileUtil.java index 1bbe17d01..5ec207467 100644 --- a/amoro-format-iceberg/src/main/java/org/apache/amoro/utils/TableFileUtil.java +++ b/amoro-format-iceberg/src/main/java/org/apache/amoro/utils/TableFileUtil.java @@ -28,6 +28,7 @@ import org.slf4j.LoggerFactory; import java.io.File; import java.net.URI; +import java.util.HashSet; import java.util.Set; import java.util.concurrent.ExecutorService; import java.util.concurrent.atomic.AtomicInteger; @@ -75,15 +76,25 @@ public class TableFileUtil { return filePath.substring(0, lastSlash); } + public static void deleteEmptyDirectory( + AuthenticatedFileIO io, String directoryPath, Set<String> exclude) { + deleteEmptyDirectory(io, directoryPath, exclude, new HashSet<>()); + } + /** * Try to recursiveDelete the empty directory * * @param io mixed-format file io * @param directoryPath directory location * @param exclude the directory will not be deleted + * @param directoriesToBeDeleted: all the directories that need to be deleted directories that + * need to be deleted */ public static void deleteEmptyDirectory( - AuthenticatedFileIO io, String directoryPath, Set<String> exclude) { + AuthenticatedFileIO io, + String directoryPath, + Set<String> exclude, + Set<String> directoriesToBeDeleted) { if (directoryPath == null || directoryPath.isEmpty()) { return; } @@ -108,7 +119,10 @@ public class TableFileUtil { if (io.asFileSystemIO().isEmptyDirectory(directoryPath)) { io.asFileSystemIO().deletePrefix(directoryPath); LOG.debug("success delete empty directory {}", directoryPath); - deleteEmptyDirectory(io, parent, exclude); + // for parent must be deleted after the sub-directory + if (directoriesToBeDeleted == null || !directoriesToBeDeleted.contains(parent)) { + deleteEmptyDirectory(io, parent, exclude); + } } } diff --git a/amoro-ams/src/test/java/org/apache/amoro/server/util/TestRollingFileCleaner.java b/amoro-format-iceberg/src/test/java/org/apache/amoro/formats/iceberg/utils/TestRollingFileCleaner.java similarity index 95% rename from amoro-ams/src/test/java/org/apache/amoro/server/util/TestRollingFileCleaner.java rename to amoro-format-iceberg/src/test/java/org/apache/amoro/formats/iceberg/utils/TestRollingFileCleaner.java index 5931f72f3..e3461b588 100644 --- a/amoro-ams/src/test/java/org/apache/amoro/server/util/TestRollingFileCleaner.java +++ b/amoro-format-iceberg/src/test/java/org/apache/amoro/formats/iceberg/utils/TestRollingFileCleaner.java @@ -16,11 +16,10 @@ * limitations under the License. */ -package org.apache.amoro.server.util; +package org.apache.amoro.formats.iceberg.utils; import org.apache.amoro.io.AuthenticatedFileIO; import org.apache.amoro.io.AuthenticatedFileIOAdapter; -import org.apache.amoro.server.utils.RollingFileCleaner; import org.apache.amoro.shade.guava32.com.google.common.collect.Sets; import org.apache.iceberg.inmemory.InMemoryFileIO; import org.junit.jupiter.api.Assertions;
