Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions include/paimon/orphan_files_cleaner.h
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,11 @@ class PAIMON_EXPORT OrphanFilesCleaner {
/// files.
virtual Result<std::set<std::string>> Clean() = 0;

/// Retrieve metrics related to orphan files cleaning operations.
///
/// @return A shared pointer to a `Metrics` object containing cleaning metrics.
virtual std::shared_ptr<Metrics> GetMetrics() const = 0;

protected:
OrphanFilesCleaner() = default;
};
Expand Down
23 changes: 21 additions & 2 deletions src/paimon/core/operation/orphan_files_cleaner_impl.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -25,14 +25,17 @@

#include "fmt/format.h"
#include "paimon/common/executor/future.h"
#include "paimon/common/metrics/metrics_impl.h"
#include "paimon/common/utils/path_util.h"
#include "paimon/common/utils/scope_guard.h"
#include "paimon/common/utils/string_utils.h"
#include "paimon/core/manifest/manifest_entry.h"
#include "paimon/core/manifest/manifest_file.h"
#include "paimon/core/manifest/manifest_file_meta.h"
#include "paimon/core/manifest/manifest_list.h"
#include "paimon/core/operation/metrics/clean_metrics.h"
#include "paimon/core/snapshot.h"
#include "paimon/core/utils/duration.h"
#include "paimon/core/utils/file_store_path_factory.h"
#include "paimon/core/utils/snapshot_manager.h"
#include "paimon/status.h"
Expand Down Expand Up @@ -61,7 +64,8 @@ OrphanFilesCleanerImpl::OrphanFilesCleanerImpl(
manifest_file_(manifest_file),
manifest_list_(manifest_list),
older_than_ms_(older_than_ms),
should_be_retained_(should_be_retained) {}
should_be_retained_(should_be_retained),
metrics_(std::make_shared<MetricsImpl>()) {}

bool OrphanFilesCleanerImpl::SupportToClean(const std::string& file_name) {
static std::vector<std::pair<std::string, std::string>> supported_pattern = {
Expand Down Expand Up @@ -91,16 +95,18 @@ Result<std::set<std::string>> OrphanFilesCleanerImpl::Clean() {
"OrphanFilesCleaner do not support cleaning table with branch");
}
PAIMON_ASSIGN_OR_RAISE(std::set<std::string> all_dirs, ListPaimonFileDirs());
PAIMON_ASSIGN_OR_RAISE(std::set<std::string> used_file_names, GetUsedFiles());
Duration duration;
std::vector<std::future<std::vector<std::unique_ptr<FileStatus>>>> file_statuses_futures;
for (const auto& dir : all_dirs) {
file_statuses_futures.push_back(
Via(executor_.get(), [this, dir] { return TryBestListingDirs(dir); }));
}
PAIMON_ASSIGN_OR_RAISE(std::set<std::string> used_file_names, GetUsedFiles());

std::set<std::string> need_to_deletes;
std::vector<std::future<void>> futures;
ScopeGuard guard([&futures]() { Wait(futures); });
uint64_t file_statuses_duration = duration.Reset();
for (const auto& file_statuses : CollectAll(file_statuses_futures)) {
for (const auto& file_status : file_statuses) {
if (file_status->IsDir()) {
Expand Down Expand Up @@ -129,10 +135,17 @@ Result<std::set<std::string>> OrphanFilesCleanerImpl::Clean() {
}
}
}
metrics_->SetCounter(CleanMetrics::CLEAN_SCAN_ORPHAN_FILES_DURATION, duration.Get());
metrics_->SetCounter(CleanMetrics::CLEAN_LIST_FILE_STATUS_DURATION, file_statuses_duration);
metrics_->SetCounter(CleanMetrics::CLEAN_LIST_FILE_STATUS_TASKS,
static_cast<uint64_t>(file_statuses_futures.size()));
metrics_->SetCounter(CleanMetrics::CLEAN_ORPHAN_FILES,
static_cast<uint64_t>(need_to_deletes.size()));
return need_to_deletes;
}

Result<std::set<std::string>> OrphanFilesCleanerImpl::ListPaimonFileDirs() const {
Duration duration;
std::set<std::string> paimon_file_dirs;
paimon_file_dirs.insert(snapshot_manager_->SnapshotDirectory());
paimon_file_dirs.insert(FileStorePathFactory::ManifestPath(root_path_));
Expand All @@ -156,6 +169,9 @@ Result<std::set<std::string>> OrphanFilesCleanerImpl::ListPaimonFileDirs() const
// ListFileDirs(external_path, partition_keys_.size());
// paimon_file_dirs.insert(external_file_dirs.begin(), external_file_dirs.end());
// }
metrics_->SetCounter(CleanMetrics::CLEAN_LIST_DIRECTORIES_DURATION, duration.Get());
metrics_->SetCounter(CleanMetrics::CLEAN_LIST_DIRECTORIES,
static_cast<uint64_t>(paimon_file_dirs.size()));
return paimon_file_dirs;
}

Expand Down Expand Up @@ -225,6 +241,7 @@ Result<std::set<std::string>> OrphanFilesCleanerImpl::GetUsedFiles() const {
// TODO(jinli.zjw): consider changelog(add tests), stats
used_files.insert(SnapshotManager::EARLIEST);
used_files.insert(SnapshotManager::LATEST);
Duration duration;
PAIMON_ASSIGN_OR_RAISE(std::vector<Snapshot> snapshots, snapshot_manager_->GetAllSnapshots());
for (const auto& snapshot : snapshots) {
used_files.insert(SnapshotManager::SNAPSHOT_PREFIX + std::to_string(snapshot.Id()));
Expand Down Expand Up @@ -257,6 +274,8 @@ Result<std::set<std::string>> OrphanFilesCleanerImpl::GetUsedFiles() const {
}
}
}
metrics_->SetCounter(CleanMetrics::CLEAN_LIST_USED_FILES_DURATION, duration.Get());
metrics_->SetCounter(CleanMetrics::CLEAN_USED_FILES, static_cast<uint64_t>(used_files.size()));
return used_files;
}

Expand Down
7 changes: 7 additions & 0 deletions src/paimon/core/operation/orphan_files_cleaner_impl.h
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
#include "paimon/core/snapshot.h"
#include "paimon/fs/file_system.h"
#include "paimon/memory/memory_pool.h"
#include "paimon/metrics.h"
#include "paimon/orphan_files_cleaner.h"
#include "paimon/result.h"

Expand Down Expand Up @@ -62,6 +63,10 @@ class OrphanFilesCleanerImpl : public OrphanFilesCleaner {

Result<std::set<std::string>> Clean() override;

std::shared_ptr<Metrics> GetMetrics() const override {
return metrics_;
}

private:
Result<std::set<std::string>> ListPaimonFileDirs() const;
std::vector<std::unique_ptr<FileStatus>> TryBestListingDirs(const std::string& path) const;
Expand All @@ -86,5 +91,7 @@ class OrphanFilesCleanerImpl : public OrphanFilesCleaner {
std::shared_ptr<ManifestList> manifest_list_;
int64_t older_than_ms_;
std::function<bool(const std::string&)> should_be_retained_;

std::shared_ptr<Metrics> metrics_;
};
} // namespace paimon
Loading