diff --git a/mkdocs/docs/file-io.md b/mkdocs/docs/file-io.md index 5b123addf..64b3682e7 100644 --- a/mkdocs/docs/file-io.md +++ b/mkdocs/docs/file-io.md @@ -80,6 +80,28 @@ Without credentials, the AWS default credential chain is used, which covers environment variables, the shared configuration file, and the various role and identity providers. +### Bulk deletion + +S3 `DeleteFiles` attempts each file, logs each failed deletion, and returns the +failure count. Exceptions thrown by a delete are counted as failures too. +Deletion is sequential by default. To enable concurrency, configure an +application-owned `iceberg::Executor` using +`iceberg::arrow::SetS3FileIODeleteExecutor(executor)` from +`iceberg/arrow/arrow_io_util.h`. + +The setting applies to all S3 FileIO instances, including ones created by the +registry or a REST catalog. Configure it before starting deletes. The executor +controls concurrency; each call submits at most 64 worker tasks, which share +the file list rather than creating a task and future for every file. + +`DeleteFiles` is synchronous and waits for accepted tasks even if submission +fails. If the executor rejects a worker, the rejection is logged and the +calling thread deletes the files no accepted worker has taken. Keep the +executor alive until all calls finish, then reset the setting to `nullptr` +before destroying it. Replace the executor only while no deletes are running. +The executor must make progress while callers wait; do not call `DeleteFiles` +from its workers unless it supports nested blocking work. + ### S3-compatible storage Stores that speak the S3 API are served by the same implementation. The scheme diff --git a/src/iceberg/arrow/arrow_io_internal.h b/src/iceberg/arrow/arrow_io_internal.h index a6b85b6c9..401f107d0 100644 --- a/src/iceberg/arrow/arrow_io_internal.h +++ b/src/iceberg/arrow/arrow_io_internal.h @@ -20,6 +20,7 @@ #pragma once #include +#include #include #include #include @@ -30,9 +31,15 @@ #include "iceberg/file_io.h" #include "iceberg/iceberg_bundle_export.h" +#include "iceberg/util/executor.h" namespace iceberg::arrow { +/// \brief Delete files, count failures, and wait for accepted tasks. +ICEBERG_BUNDLE_EXPORT Status +BulkDeleteFiles(const std::vector& file_locations, Executor* executor, + const std::function& delete_file); + /// \brief Open a FileIO input as an Arrow input stream. /// /// Uses ArrowFileSystemFileIO's native Arrow stream directly when possible and falls diff --git a/src/iceberg/arrow/arrow_io_util.h b/src/iceberg/arrow/arrow_io_util.h index 925ee12a2..f5e4170fe 100644 --- a/src/iceberg/arrow/arrow_io_util.h +++ b/src/iceberg/arrow/arrow_io_util.h @@ -29,6 +29,7 @@ #include "iceberg/file_io.h" #include "iceberg/iceberg_bundle_export.h" #include "iceberg/result.h" +#include "iceberg/type_fwd.h" namespace iceberg::arrow { @@ -47,6 +48,15 @@ ICEBERG_BUNDLE_EXPORT std::unique_ptr MakeLocalFileIO(); ICEBERG_BUNDLE_EXPORT Result> MakeS3FileIO( const std::unordered_map& properties = {}); +/// \brief Set the shared executor for all S3 bulk deletes. +/// +/// Includes registry-created FileIOs; nullptr (default) runs sequentially. +/// The caller owns the executor. Configure only while no deletes are active, +/// and clear it before destruction. DeleteFiles waits for accepted work on error, +/// and deletes on the calling thread what a rejected worker would have. +/// See Executor for blocking requirements. +ICEBERG_BUNDLE_EXPORT void SetS3FileIODeleteExecutor(Executor* executor); + /// \brief Finalize (clean up) the Arrow S3 subsystem. /// /// Must be called before process exit if S3 was initialized, otherwise Arrow's diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index 4fbf33098..2028a254a 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -18,7 +18,11 @@ */ #include +#include #include +#include +#include +#include #include #include #include @@ -39,12 +43,118 @@ #include "iceberg/arrow/arrow_status_internal.h" #include "iceberg/arrow/s3/s3_properties.h" #include "iceberg/logging/log_macros.h" +#include "iceberg/logging/logger.h" +#include "iceberg/util/executor.h" #include "iceberg/util/macros.h" #include "iceberg/util/property_util.h" #include "iceberg/util/string_util.h" namespace iceberg::arrow { +namespace { + +std::atomic delete_executor{nullptr}; + +// Bounds the tasks and futures of one call, however many files it deletes. +constexpr size_t kMaxDeleteWorkers = 64; + +struct DeleteTasks { + std::vector> futures; + + ~DeleteTasks() { + for (auto& future : futures) { + if (future.valid()) { + future.wait(); + } + } + } +}; + +// Submit may throw instead of returning an error; both reject the worker. +Status SubmitWorker(Executor& executor, ExecutorTask task) { + try { + return executor.Submit(std::move(task)); + } catch (const std::exception& e) { + return IOError("Submit threw an exception: {}", e.what()); + } catch (...) { + return IOError("Submit threw an unknown exception"); + } +} + +} // namespace + +void SetS3FileIODeleteExecutor(Executor* executor) { + delete_executor.store(executor, std::memory_order_release); +} + +Status BulkDeleteFiles(const std::vector& file_locations, Executor* executor, + const std::function& delete_file) { + auto logger = GetCurrentLogger(); + std::atomic next = 0; + std::atomic failed = 0; + auto work = [&] { + ScopedLogger bind(logger); + while (true) { + const auto index = next.fetch_add(1, std::memory_order_relaxed); + if (index >= file_locations.size()) { + return; + } + const auto& location = file_locations[index]; + Status status; + try { + status = delete_file(location); + } catch (const std::exception& e) { + status = IOError("Delete threw an exception: {}", e.what()); + } catch (...) { + status = IOError("Delete threw an unknown exception"); + } + if (!status.has_value()) { + failed.fetch_add(1, std::memory_order_relaxed); + ICEBERG_LOG_WARN("Failed to delete {}: {}", location, status.error().message); + } + } + }; + + try { + if (executor == nullptr) { + work(); + } else { + const auto workers = std::min(kMaxDeleteWorkers, file_locations.size()); + DeleteTasks tasks; + tasks.futures.reserve(workers); + size_t accepted = 0; + for (; accepted < workers; ++accepted) { + std::packaged_task task(work); + // Submit may accept the task before failing. + tasks.futures.push_back(task.get_future()); + ExecutorTask executor_task([task = std::move(task)]() mutable { task(); }); + if (auto status = SubmitWorker(*executor, std::move(executor_task)); + !status.has_value()) { + // Take what the accepted workers have not, so every file is still + // attempted and counted even if none was accepted. + ICEBERG_LOG_WARN("Delete worker rejected; deleting on the calling thread: {}", + status.error().message); + work(); + break; + } + } + // A rejected worker's future is only waited for, by `tasks`. + for (size_t i = 0; i < accepted; ++i) { + tasks.futures[i].get(); + } + } + } catch (const std::exception& e) { + return IOError("Bulk delete execution failed: {}", e.what()); + } catch (...) { + return IOError("Bulk delete execution failed with an unknown exception"); + } + if (failed > 0) { + return IOError("Failed to delete {} of {} files", failed.load(), + file_locations.size()); + } + return {}; +} + #if ICEBERG_S3_ENABLED namespace { @@ -338,8 +448,7 @@ Status ArrowS3FileIO::DeleteFile(const std::string& file_location) { } Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations) { - // One snapshot so the whole batch matches the same delegate generation; only - // ever a handful of delegates, so a linear scan beats hashing. + // Keep one delegate snapshot for the whole batch. std::shared_ptr fallback; DelegatesByPrefix by_prefix; { @@ -347,23 +456,11 @@ Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations fallback = default_file_io_; by_prefix = file_io_by_prefix_; } - std::vector, std::vector>> - locations_by_io; - for (const auto& file_location : file_locations) { - auto file_io = MatchDelegate(fallback, by_prefix, file_location); - auto it = std::ranges::find_if( - locations_by_io, [&](const auto& entry) { return entry.first == file_io; }); - if (it == locations_by_io.end()) { - locations_by_io.emplace_back(std::move(file_io), - std::vector{file_location}); - } else { - it->second.push_back(file_location); - } - } - for (auto& [file_io, locations] : locations_by_io) { - ICEBERG_RETURN_UNEXPECTED(file_io->DeleteFiles(locations)); - } - return {}; + return BulkDeleteFiles( + file_locations, delete_executor.load(std::memory_order_acquire), + [&](const std::string& location) { + return MatchDelegate(fallback, by_prefix, location)->DeleteFile(location); + }); } } // namespace diff --git a/src/iceberg/test/arrow_s3_file_io_test.cc b/src/iceberg/test/arrow_s3_file_io_test.cc index 24628949a..19d66fe97 100644 --- a/src/iceberg/test/arrow_s3_file_io_test.cc +++ b/src/iceberg/test/arrow_s3_file_io_test.cc @@ -19,11 +19,14 @@ #include #include +#include #include +#include #include #include #include #include +#include #include #include #include @@ -38,12 +41,16 @@ # include #endif +#include "iceberg/arrow/arrow_io_internal.h" #include "iceberg/arrow/arrow_io_util.h" +#include "iceberg/arrow/arrow_register.h" #include "iceberg/arrow/s3/s3_properties.h" #include "iceberg/file_io.h" +#include "iceberg/file_io_registry.h" #include "iceberg/logging/logger.h" #include "iceberg/result.h" #include "iceberg/storage_credential.h" +#include "iceberg/test/executor.h" #include "iceberg/test/logging_test_helpers.h" #include "iceberg/test/matchers.h" #include "iceberg/util/macros.h" @@ -156,8 +163,219 @@ Status CheckReadWrite(FileIO& io, const std::string& object_uri, return io.DeleteFile(object_uri); } +class InlineDeleteExecutor : public Executor { + public: + Status Submit(ExecutorTask task) override { + ++submissions; + std::move(task)(); + return {}; + } + + int submissions = 0; +}; + +class ScopedDeleteExecutor { + public: + explicit ScopedDeleteExecutor(Executor* executor) { + SetS3FileIODeleteExecutor(executor); + } + ~ScopedDeleteExecutor() { SetS3FileIODeleteExecutor(nullptr); } +}; + +enum class SubmitFailure { kStatus, kException, kAcceptedException }; + +class FailingDeleteExecutor : public Executor { + public: + FailingDeleteExecutor(SubmitFailure failure, std::shared_future release) + : failure_(failure), release_(std::move(release)) {} + + ~FailingDeleteExecutor() override { Wait(); } + + void Wait() { + for (auto& thread : threads_) { + if (thread.joinable()) { + thread.join(); + } + } + } + + Status Submit(ExecutorTask task) override { + ++submissions_; + if (submissions_ == 1 || failure_ == SubmitFailure::kAcceptedException) { + const bool first = submissions_ == 1; + const bool wait = !first || failure_ != SubmitFailure::kAcceptedException; + threads_.emplace_back([this, first, wait, task = std::move(task)]() mutable { + if (wait) { + release_.wait(); + } + std::move(task)(); + ++completed; + if (first) { + first_completed.set_value(); + } + }); + } + if (submissions_ == 2) { + submission_failed.set_value(); + if (failure_ == SubmitFailure::kStatus) { + return ServiceUnavailable("executor rejected deletion"); + } + throw std::runtime_error("executor rejected deletion"); + } + return {}; + } + + std::promise submission_failed; + std::promise first_completed; + std::atomic completed = 0; + + private: + SubmitFailure failure_; + std::shared_future release_; + int submissions_ = 0; + std::vector threads_; +}; + } // namespace +TEST(BulkDeleteFilesTest, SequentialDeletionRunsOnTheCaller) { + const std::vector paths = {"first", "second", "third"}; + std::vector attempted; + const auto caller = std::this_thread::get_id(); + EXPECT_THAT(BulkDeleteFiles(paths, nullptr, + [&](const std::string& path) -> Status { + EXPECT_EQ(std::this_thread::get_id(), caller); + attempted.push_back(path); + return {}; + }), + IsOk()); + EXPECT_EQ(attempted, paths); +} + +TEST(BulkDeleteFilesTest, CountsThrownDeletesAndAttemptsRemainingFiles) { + for (bool parallel : {false, true}) { + SCOPED_TRACE(parallel); + test::ThreadExecutor executor; + const std::vector paths = {"status", "exception", "unknown", "ok"}; + std::atomic attempted = 0; + auto logger = std::make_shared(); + ScopedLogger bind(logger); + auto status = BulkDeleteFiles(paths, parallel ? &executor : nullptr, + [&](const std::string& path) -> Status { + ++attempted; + if (path == "status") { + return IOError("delete failed"); + } + if (path == "exception") { + throw std::runtime_error("delete threw"); + } + if (path == "unknown") { + throw 42; + } + return {}; + }); + EXPECT_THAT(status, HasErrorMessage("Failed to delete 3 of 4 files")); + EXPECT_EQ(attempted, 4); + EXPECT_EQ(logger->count(), 3); + } +} + +TEST(BulkDeleteFilesTest, DrainsAcceptedTasksWhenSubmissionFails) { + for (auto failure : {SubmitFailure::kStatus, SubmitFailure::kException, + SubmitFailure::kAcceptedException}) { + SCOPED_TRACE(static_cast(failure)); + std::promise release; + FailingDeleteExecutor executor(failure, release.get_future().share()); + auto submission_failed = executor.submission_failed.get_future(); + auto first_completed = executor.first_completed.get_future(); + const std::vector paths = {"first", "second", "third"}; + std::atomic attempted = 0; + auto logger = std::make_shared(); + auto run = std::async(std::launch::async, [&] { + ScopedLogger bind(logger); + return BulkDeleteFiles(paths, &executor, [&](const std::string&) -> Status { + ++attempted; + return {}; + }); + }); + const auto submitted = submission_failed.wait_for(std::chrono::seconds(5)); + auto first = std::future_status::ready; + if (failure == SubmitFailure::kAcceptedException) { + // Leave only the last task blocked. + first = first_completed.wait_for(std::chrono::seconds(5)); + } + const auto returned = run.wait_for(std::chrono::milliseconds(100)); + // Unblock workers before assertions. + release.set_value(); + auto status = run.get(); + executor.Wait(); + EXPECT_EQ(submitted, std::future_status::ready); + EXPECT_EQ(first, std::future_status::ready); + EXPECT_EQ(returned, std::future_status::timeout); + // The caller deletes what the accepted workers did not; the rejection is + // logged, not returned. + EXPECT_THAT(status, IsOk()); + EXPECT_EQ(attempted, 3); + const auto records = logger->records(); + ASSERT_EQ(records.size(), 1); + EXPECT_THAT(records[0].message, ::testing::HasSubstr("executor rejected deletion")); + EXPECT_EQ(executor.completed, failure == SubmitFailure::kAcceptedException ? 2 : 1); + } +} + +TEST(BulkDeleteFilesTest, DeletesOnTheCallerWhenNoWorkerIsAccepted) { + test::ThreadExecutor executor(ServiceUnavailable("executor is shut down")); + const std::vector paths = {"first", "second", "third"}; + std::vector attempted; + const auto caller = std::this_thread::get_id(); + EXPECT_THAT(BulkDeleteFiles(paths, &executor, + [&](const std::string& path) -> Status { + EXPECT_EQ(std::this_thread::get_id(), caller); + attempted.push_back(path); + if (path == "second") { + return IOError("delete failed"); + } + return {}; + }), + HasErrorMessage("Failed to delete 1 of 3 files")); + EXPECT_EQ(attempted, paths); + EXPECT_EQ(executor.submit_count(), 1); +} + +TEST(BulkDeleteFilesTest, TaskCountDoesNotGrowWithTheFileList) { + std::vector task_counts; + for (size_t size : {4096, 16384}) { + InlineDeleteExecutor executor; + std::vector paths; + std::vector attempts(size, 0); + for (size_t i = 0; i < size; ++i) { + paths.push_back(std::to_string(i)); + } + EXPECT_THAT(BulkDeleteFiles(paths, &executor, + [&](const std::string& path) -> Status { + ++attempts[std::stoul(path)]; + return {}; + }), + IsOk()); + EXPECT_THAT(attempts, ::testing::Each(1)); + EXPECT_LT(executor.submissions, size); + task_counts.push_back(executor.submissions); + } + EXPECT_EQ(task_counts[0], task_counts[1]); +} + +TEST(BulkDeleteFilesTest, EmptyDeleteDoesNotSubmitTasks) { + InlineDeleteExecutor executor; + EXPECT_THAT(BulkDeleteFiles({}, &executor, + [](const std::string&) -> Status { + ADD_FAILURE() + << "An empty delete should not invoke the callback"; + return {}; + }), + IsOk()); + EXPECT_EQ(executor.submissions, 0); +} + TEST_F(ArrowS3FileIOTest, Create) { auto result = MakeS3FileIO({}); ASSERT_THAT(result, IsOk()); @@ -299,6 +517,26 @@ TEST_F(ArrowS3FileIOTest, RejectsIncompleteStaticCredentials) { "S3 client access key ID and secret access key must be set")); } +TEST_F(ArrowS3FileIOTest, SharesDeleteExecutorWithRegistryCreatedFileIOs) { + RegisterAll(); + ICEBERG_UNWRAP_OR_FAIL(auto direct, MakeS3FileIO({})); + ICEBERG_UNWRAP_OR_FAIL(auto registered, + FileIORegistry::Load(FileIORegistry::kArrowS3FileIO, {})); + InlineDeleteExecutor executor; + ScopedDeleteExecutor configured(&executor); + // Invalid schemes avoid network I/O. + for (auto* io : {direct.get(), registered.get()}) { + EXPECT_THAT(io->DeleteFiles({"invalid://bucket/file"}), + HasErrorMessage("Failed to delete 1 of 1 files")); + } + EXPECT_EQ(executor.submissions, 2); + + SetS3FileIODeleteExecutor(nullptr); + EXPECT_THAT(direct->DeleteFiles({"invalid://bucket/file"}), + HasErrorMessage("Failed to delete 1 of 1 files")); + EXPECT_EQ(executor.submissions, 2); +} + TEST_F(ArrowS3FileIOTest, ReadWrite) { if (!HasIntegrationEnv()) { GTEST_SKIP() << "Set ICEBERG_TEST_S3_URI to enable S3 IO test"; @@ -353,6 +591,56 @@ TEST_F(ArrowS3FileIOTest, LongestCredentialPrefix) { IsOk()); } +TEST_F(ArrowS3FileIOTest, DeleteFilesAttemptsEveryFile) { + if (!HasIntegrationEnv()) { + GTEST_SKIP() << "Set ICEBERG_TEST_S3_URI to enable S3 IO test"; + } + + auto properties = PropertiesFromEnv(); + if (properties.empty()) { + GTEST_SKIP() << "Set S3 properties to enable credential routing test"; + } + test::ThreadExecutor executor; + ScopedDeleteExecutor configured(&executor); + + auto io_res = MakeS3FileIO(properties); + ASSERT_THAT(io_res, IsOk()); + auto io = std::move(io_res).value(); + auto* credentialed = io->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + const auto denied = ObjectUri("delete_denied/"); + const auto allowed = ObjectUri("delete_allowed/") + "only"; + const std::vector paths = {denied + "first", allowed, denied + "second"}; + for (const auto& path : paths) { + ASSERT_THAT(io->WriteFile(path, "payload"), IsOk()); + } + + // Deletes under `denied` fail; the one between them must still run. + auto bad_properties = properties; + for (const auto& [key, value] : BadS3Credentials()) { + bad_properties.insert_or_assign(key, value); + } + ASSERT_THAT(credentialed->SetStorageCredentials( + {{.prefix = denied, .config = std::move(bad_properties)}}), + IsOk()); + // Each failure reaches the caller's logger, whichever thread hit it. + auto logger = std::make_shared(); + { + ScopedLogger bind(logger); + EXPECT_THAT(io->DeleteFiles(paths), HasErrorMessage("Failed to delete 2 of 3 files")); + } + EXPECT_EQ(std::ranges::count_if( + logger->records(), + [](const LogMessage& record) { return record.level == LogLevel::kWarn; }), + 2); + EXPECT_FALSE(io->ReadFile(allowed, std::nullopt).has_value()); + + // The denied files remain: Arrow fails to delete a missing object. + ASSERT_THAT(credentialed->SetStorageCredentials({}), IsOk()); + EXPECT_THAT(io->DeleteFiles({paths[0], paths[2]}), IsOk()); +} + // The credential is vended under the oss spelling and the object addressed as // `s3://`, so they only meet through canonicalization — and every other path // to authentication is broken. (rest_arrow_file_io_test covers the mirrored diff --git a/src/iceberg/util/executor.h b/src/iceberg/util/executor.h index 749502b17..625c833f1 100644 --- a/src/iceberg/util/executor.h +++ b/src/iceberg/util/executor.h @@ -33,7 +33,7 @@ namespace iceberg { using ExecutorTask = FnOnce; -/// \brief Schedules iceberg-cpp internal planning tasks. +/// \brief Schedules iceberg-cpp tasks. /// /// Public APIs that accept an executor remain synchronous: the calling thread may block /// while waiting for submitted tasks to finish. Callers must ensure the executor can @@ -41,8 +41,8 @@ using ExecutorTask = FnOnce; /// the same bounded executor's worker threads can deadlock unless the executor supports /// nested blocking work. /// -/// When an executor is configured, planning callbacks may be called concurrently. Any -/// shared mutable state captured by those callbacks must be synchronized by the caller. +/// When an executor is configured, callbacks may be called concurrently. Any shared +/// mutable state captured by those callbacks must be synchronized by the caller. class ICEBERG_EXPORT Executor { public: virtual ~Executor() = default;