From 1c0b5bb2158d71f2a07708e0dc2fec165715b27b Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Sat, 26 Sep 2026 02:22:57 -0400 Subject: [PATCH 1/8] fix(io): attempt every S3 delete and report how many failed ArrowS3FileIO::DeleteFiles grouped locations by credential prefix and returned at the first group that failed, so files in later groups were never attempted. Java's S3FileIO attempts every batch, logs each failed path and reports the failure count. Delete each file through its delegate, log each failure, and return "Failed to delete N of M files" at the end. This sends the same requests as before, since Arrow's S3 DeleteFiles also deletes one file at a time. ResolvingFileIO stays fail-fast across delegates, as it is in Java. --- src/iceberg/arrow/s3/arrow_s3_file_io.cc | 25 +++++++-------- src/iceberg/test/arrow_s3_file_io_test.cc | 39 +++++++++++++++++++++++ 2 files changed, 50 insertions(+), 14 deletions(-) diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index 4fbf330985..552f67de69 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -338,8 +338,9 @@ 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. + // Like Java's S3FileIO, keep going after a failure and report the count. + // Arrow's S3 DeleteFiles deletes one file at a time too. One snapshot so the + // whole batch matches the same delegate generation. std::shared_ptr fallback; DelegatesByPrefix by_prefix; { @@ -347,21 +348,17 @@ 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; + size_t failed = 0; 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); + if (auto status = + MatchDelegate(fallback, by_prefix, file_location)->DeleteFile(file_location); + !status.has_value()) { + ICEBERG_LOG_WARN("Failed to delete {}: {}", file_location, status.error().message); + ++failed; } } - for (auto& [file_io, locations] : locations_by_io) { - ICEBERG_RETURN_UNEXPECTED(file_io->DeleteFiles(locations)); + if (failed > 0) { + return IOError("Failed to delete {} of {} files", failed, file_locations.size()); } return {}; } diff --git a/src/iceberg/test/arrow_s3_file_io_test.cc b/src/iceberg/test/arrow_s3_file_io_test.cc index 24628949a9..d8c0150614 100644 --- a/src/iceberg/test/arrow_s3_file_io_test.cc +++ b/src/iceberg/test/arrow_s3_file_io_test.cc @@ -353,6 +353,45 @@ 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"; + } + + 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()); + EXPECT_THAT(io->DeleteFiles(paths), HasErrorMessage("Failed to delete 2 of 3 files")); + 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 From bef73edc86aa993a2f2016d87ac6195e4d207991 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Sat, 26 Sep 2026 06:45:25 -0400 Subject: [PATCH 2/8] perf(io): run S3 deletes on s3.delete.num-threads threads Java's S3FileIO runs its deletes on an executor sized by s3.delete.num-threads, which defaults to the number of processors. DeleteFiles now does the same, with the calling thread taking part. Java also packs keys into DeleteObjects batches. Arrow has no public batch delete for S3, so each thread still deletes one file at a time, as Arrow's own DeleteFiles does. --- src/iceberg/arrow/s3/arrow_s3_file_io.cc | 54 +++++++++++++++++------ src/iceberg/arrow/s3/s3_properties.h | 2 + src/iceberg/test/arrow_s3_file_io_test.cc | 11 +++++ 3 files changed, 54 insertions(+), 13 deletions(-) diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index 552f67de69..87011d7f0e 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -18,6 +18,7 @@ */ #include +#include #include #include #include @@ -25,6 +26,7 @@ #include #include #include +#include #include #include #include @@ -181,9 +183,11 @@ std::string CanonicalizeS3Scheme(std::string_view location) { class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { public: ArrowS3FileIO(std::shared_ptr<::arrow::fs::FileSystem> arrow_fs, - std::unordered_map default_properties) + std::unordered_map default_properties, + size_t delete_threads) : default_file_io_(std::make_shared(std::move(arrow_fs))), - default_properties_(std::move(default_properties)) {} + default_properties_(std::move(default_properties)), + delete_threads_(delete_threads) {} Result> NewInputFile(std::string file_location) override; @@ -237,6 +241,7 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { std::shared_ptr default_file_io_; std::unordered_map default_properties_; + size_t delete_threads_; // Guards everything below; shared because reads happen per file operation. mutable std::shared_mutex mutex_; std::vector storage_credentials_; @@ -338,8 +343,8 @@ Status ArrowS3FileIO::DeleteFile(const std::string& file_location) { } Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations) { - // Like Java's S3FileIO, keep going after a failure and report the count. - // Arrow's S3 DeleteFiles deletes one file at a time too. One snapshot so the + // Like Java's S3FileIO: delete concurrently, keep going after a failure and + // report the count. Arrow has no batch delete for S3. One snapshot so the // whole batch matches the same delegate generation. std::shared_ptr fallback; DelegatesByPrefix by_prefix; @@ -348,17 +353,31 @@ Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations fallback = default_file_io_; by_prefix = file_io_by_prefix_; } - size_t failed = 0; - for (const auto& file_location : file_locations) { - if (auto status = - MatchDelegate(fallback, by_prefix, file_location)->DeleteFile(file_location); - !status.has_value()) { - ICEBERG_LOG_WARN("Failed to delete {}: {}", file_location, status.error().message); - ++failed; + std::atomic next = 0; + std::atomic failed = 0; + auto delete_remaining = [&] { + for (size_t i = next++; i < file_locations.size(); i = next++) { + const auto& file_location = file_locations[i]; + if (auto status = MatchDelegate(fallback, by_prefix, file_location) + ->DeleteFile(file_location); + !status.has_value()) { + ICEBERG_LOG_WARN("Failed to delete {}: {}", file_location, + status.error().message); + ++failed; + } } + }; + { + // Plus the calling thread. + std::vector helpers; + for (size_t i = 1; i < std::min(delete_threads_, file_locations.size()); ++i) { + helpers.emplace_back(delete_remaining); + } + delete_remaining(); } if (failed > 0) { - return IOError("Failed to delete {} of {} files", failed, file_locations.size()); + return IOError("Failed to delete {} of {} files", failed.load(), + file_locations.size()); } return {}; } @@ -369,7 +388,16 @@ Result> MakeS3FileIO( const std::unordered_map& properties) { // Uses default credentials if properties are empty. ICEBERG_ASSIGN_OR_RAISE(auto fs, BuildArrowS3FileSystem(properties)); - return std::make_unique(std::move(fs), properties); + // Java defaults to one delete thread per processor. + size_t delete_threads = std::max(1u, std::thread::hardware_concurrency()); + if (const auto* value = FindProperty(properties, S3Properties::kDeleteNumThreads); + value != nullptr) { + ICEBERG_ASSIGN_OR_RAISE(delete_threads, StringUtils::ParseNumber(*value)); + if (delete_threads == 0) { + return InvalidArgument(R"("{}" must be positive)", S3Properties::kDeleteNumThreads); + } + } + return std::make_unique(std::move(fs), properties, delete_threads); } Status FinalizeS3() { diff --git a/src/iceberg/arrow/s3/s3_properties.h b/src/iceberg/arrow/s3/s3_properties.h index 50dafa56ff..bcc1865838 100644 --- a/src/iceberg/arrow/s3/s3_properties.h +++ b/src/iceberg/arrow/s3/s3_properties.h @@ -54,6 +54,8 @@ struct S3Properties { static constexpr std::string_view kConnectTimeoutMs = "s3.connect-timeout-ms"; /// Socket timeout in milliseconds static constexpr std::string_view kSocketTimeoutMs = "s3.socket-timeout-ms"; + /// Threads DeleteFiles uses; defaults to the hardware thread count, as in Java + static constexpr std::string_view kDeleteNumThreads = "s3.delete.num-threads"; }; /// \brief URI schemes served by the Arrow S3 FileIO, lower-case. diff --git a/src/iceberg/test/arrow_s3_file_io_test.cc b/src/iceberg/test/arrow_s3_file_io_test.cc index d8c0150614..46e181d060 100644 --- a/src/iceberg/test/arrow_s3_file_io_test.cc +++ b/src/iceberg/test/arrow_s3_file_io_test.cc @@ -299,6 +299,15 @@ TEST_F(ArrowS3FileIOTest, RejectsIncompleteStaticCredentials) { "S3 client access key ID and secret access key must be set")); } +TEST_F(ArrowS3FileIOTest, RejectsInvalidDeleteThreads) { + for (std::string_view threads : {"0", "-1", "many"}) { + SCOPED_TRACE(threads); + EXPECT_THAT(MakeS3FileIO({{std::string(S3Properties::kDeleteNumThreads), + std::string(threads)}}), + IsError(ErrorKind::kInvalidArgument)); + } +} + TEST_F(ArrowS3FileIOTest, ReadWrite) { if (!HasIntegrationEnv()) { GTEST_SKIP() << "Set ICEBERG_TEST_S3_URI to enable S3 IO test"; @@ -362,6 +371,8 @@ TEST_F(ArrowS3FileIOTest, DeleteFilesAttemptsEveryFile) { if (properties.empty()) { GTEST_SKIP() << "Set S3 properties to enable credential routing test"; } + // A thread per file, so the deletes run concurrently. + properties[std::string(S3Properties::kDeleteNumThreads)] = "3"; auto io_res = MakeS3FileIO(properties); ASSERT_THAT(io_res, IsOk()); From aee06c29cebe9e03488e35490f6e877bf538ff80 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Sat, 26 Sep 2026 08:20:31 -0400 Subject: [PATCH 3/8] fix(io): run S3 delete helpers via std::async with the caller's logger std::jthread needs -fexperimental-library with libc++ 18 and 19, which the Clang 18+ requirement covers, so use std::async futures instead; they also wait for their thread if an exception unwinds. Helper threads now bind the caller's logger, as logger.h prescribes for thread pools. Without it their warnings went to the global logger while the returned error only carries the failure count. --- src/iceberg/arrow/s3/arrow_s3_file_io.cc | 21 ++++++++++++++------- src/iceberg/test/arrow_s3_file_io_test.cc | 11 ++++++++++- 2 files changed, 24 insertions(+), 8 deletions(-) diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index 87011d7f0e..ba44715d39 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -20,6 +20,7 @@ #include #include #include +#include #include #include #include @@ -41,6 +42,7 @@ #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/macros.h" #include "iceberg/util/property_util.h" #include "iceberg/util/string_util.h" @@ -367,13 +369,18 @@ Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations } } }; - { - // Plus the calling thread. - std::vector helpers; - for (size_t i = 1; i < std::min(delete_threads_, file_locations.size()); ++i) { - helpers.emplace_back(delete_remaining); - } - delete_remaining(); + // Plus the calling thread. Helpers log where the caller does. + auto logger = GetCurrentLogger(); + std::vector> helpers; + for (size_t i = 1; i < std::min(delete_threads_, file_locations.size()); ++i) { + helpers.push_back(std::async(std::launch::async, [&] { + ScopedLogger bind(logger); + delete_remaining(); + })); + } + delete_remaining(); + for (auto& helper : helpers) { + helper.wait(); } if (failed > 0) { return IOError("Failed to delete {} of {} files", failed.load(), diff --git a/src/iceberg/test/arrow_s3_file_io_test.cc b/src/iceberg/test/arrow_s3_file_io_test.cc index 46e181d060..f75958c70c 100644 --- a/src/iceberg/test/arrow_s3_file_io_test.cc +++ b/src/iceberg/test/arrow_s3_file_io_test.cc @@ -395,7 +395,16 @@ TEST_F(ArrowS3FileIOTest, DeleteFilesAttemptsEveryFile) { ASSERT_THAT(credentialed->SetStorageCredentials( {{.prefix = denied, .config = std::move(bad_properties)}}), IsOk()); - EXPECT_THAT(io->DeleteFiles(paths), HasErrorMessage("Failed to delete 2 of 3 files")); + // 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. From 492c6cb34751110ff95da57c6a19a71b91b24fe1 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Sat, 26 Sep 2026 08:28:47 -0400 Subject: [PATCH 4/8] fix(io): rethrow exceptions from S3 delete helpers wait() leaves an exception stored in a helper's future, so a delete that threw was neither counted nor reported, and DeleteFiles could return success. get() rethrows it on the calling thread, as the sequential loop did. --- src/iceberg/arrow/s3/arrow_s3_file_io.cc | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index ba44715d39..00afb46d97 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -380,7 +380,7 @@ Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations } delete_remaining(); for (auto& helper : helpers) { - helper.wait(); + helper.get(); // Rethrows what a helper threw. } if (failed > 0) { return IOError("Failed to delete {} of {} files", failed.load(), From dd0da21ec72e7e1ff40e49a7c0aa38b90fce7cf7 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Tue, 29 Sep 2026 02:13:24 -0400 Subject: [PATCH 5/8] docs(io): document s3.delete.num-threads --- mkdocs/docs/file-io.md | 1 + 1 file changed, 1 insertion(+) diff --git a/mkdocs/docs/file-io.md b/mkdocs/docs/file-io.md index 5b123addf7..6a5842ca7a 100644 --- a/mkdocs/docs/file-io.md +++ b/mkdocs/docs/file-io.md @@ -66,6 +66,7 @@ each file location's scheme. | `client.region` | `us-east-1` | Region to sign requests for | | `s3.endpoint` | `https://127.0.0.1:9000` | Endpoint to use instead of the AWS one. When absent, the `AWS_ENDPOINT_URL_S3` / `AWS_ENDPOINT_URL` environment variables are consulted | | `s3.path-style-access` | `true` | Address buckets as a path (`endpoint/bucket`) instead of a virtual host (`bucket.endpoint`). Only takes effect together with a custom endpoint | +| `s3.delete.num-threads` | `8` | Number of threads `DeleteFiles` deletes with. Defaults to the number of hardware threads | The following keys are specific to iceberg-cpp; they are not part of the Java Iceberg or REST specification property set: From 75adea35eee3859d84a5084d4413b8d6e23e9555 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Wed, 7 Oct 2026 22:56:01 -0400 Subject: [PATCH 6/8] refactor(io): run S3 deletes on a pool the FileIO owns Each DeleteFiles call started its own threads, so s3.delete.num-threads bounded one call rather than the FileIO: concurrent calls could run many times that number. The FileIO now owns an Arrow thread pool of that size, created with the FileIO but starting workers on the first delete, and runs each batch on it through TaskGroup, as Java's S3FileIO does with its executor. Failures are still logged per file and reported as a count. --- mkdocs/docs/file-io.md | 2 +- src/iceberg/arrow/s3/arrow_s3_file_io.cc | 81 ++++++++++++++---------- src/iceberg/arrow/s3/s3_properties.h | 3 +- 3 files changed, 52 insertions(+), 34 deletions(-) diff --git a/mkdocs/docs/file-io.md b/mkdocs/docs/file-io.md index 6a5842ca7a..8064c8ae71 100644 --- a/mkdocs/docs/file-io.md +++ b/mkdocs/docs/file-io.md @@ -66,7 +66,7 @@ each file location's scheme. | `client.region` | `us-east-1` | Region to sign requests for | | `s3.endpoint` | `https://127.0.0.1:9000` | Endpoint to use instead of the AWS one. When absent, the `AWS_ENDPOINT_URL_S3` / `AWS_ENDPOINT_URL` environment variables are consulted | | `s3.path-style-access` | `true` | Address buckets as a path (`endpoint/bucket`) instead of a virtual host (`bucket.endpoint`). Only takes effect together with a custom endpoint | -| `s3.delete.num-threads` | `8` | Number of threads `DeleteFiles` deletes with. Defaults to the number of hardware threads | +| `s3.delete.num-threads` | `8` | Size of the thread pool `DeleteFiles` runs on, shared by every call on the FileIO. Defaults to the number of hardware threads | The following keys are specific to iceberg-cpp; they are not part of the Java Iceberg or REST specification property set: diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index 00afb46d97..a0d18b776c 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -20,7 +20,7 @@ #include #include #include -#include +#include #include #include #include @@ -35,6 +35,7 @@ #include #if ICEBERG_S3_ENABLED # include +# include #endif #include "iceberg/arrow/arrow_io_internal.h" @@ -46,6 +47,7 @@ #include "iceberg/util/macros.h" #include "iceberg/util/property_util.h" #include "iceberg/util/string_util.h" +#include "iceberg/util/task_group.h" namespace iceberg::arrow { @@ -182,14 +184,29 @@ std::string CanonicalizeS3Scheme(std::string_view location) { return std::string(location); } +/// \brief Runs TaskGroup tasks on an Arrow thread pool. +class ThreadPoolExecutor final : public Executor { + public: + explicit ThreadPoolExecutor(std::shared_ptr<::arrow::internal::ThreadPool> pool) + : pool_(std::move(pool)) {} + + Status Submit(ExecutorTask task) override { + ICEBERG_ARROW_RETURN_NOT_OK(pool_->Spawn(std::move(task))); + return {}; + } + + private: + std::shared_ptr<::arrow::internal::ThreadPool> pool_; +}; + class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { public: ArrowS3FileIO(std::shared_ptr<::arrow::fs::FileSystem> arrow_fs, std::unordered_map default_properties, - size_t delete_threads) + std::shared_ptr<::arrow::internal::ThreadPool> delete_pool) : default_file_io_(std::make_shared(std::move(arrow_fs))), default_properties_(std::move(default_properties)), - delete_threads_(delete_threads) {} + delete_executor_(std::move(delete_pool)) {} Result> NewInputFile(std::string file_location) override; @@ -243,7 +260,9 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { std::shared_ptr default_file_io_; std::unordered_map default_properties_; - size_t delete_threads_; + // Bounds the deletes of every DeleteFiles call on this FileIO, like the + // executor of Java's S3FileIO. + ThreadPoolExecutor delete_executor_; // Guards everything below; shared because reads happen per file operation. mutable std::shared_mutex mutex_; std::vector storage_credentials_; @@ -345,9 +364,9 @@ Status ArrowS3FileIO::DeleteFile(const std::string& file_location) { } Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations) { - // Like Java's S3FileIO: delete concurrently, keep going after a failure and - // report the count. Arrow has no batch delete for S3. One snapshot so the - // whole batch matches the same delegate generation. + // Like Java's S3FileIO: delete on the FileIO's bounded pool, keep going after + // a failure and report the count. Arrow has no batch delete for S3. One + // snapshot so the whole batch matches the same delegate generation. std::shared_ptr fallback; DelegatesByPrefix by_prefix; { @@ -355,33 +374,25 @@ Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations fallback = default_file_io_; by_prefix = file_io_by_prefix_; } - std::atomic next = 0; + // Each failure is logged where the caller logs and only counted, so the + // returned error stays one line however large the batch. + auto logger = GetCurrentLogger(); std::atomic failed = 0; - auto delete_remaining = [&] { - for (size_t i = next++; i < file_locations.size(); i = next++) { - const auto& file_location = file_locations[i]; - if (auto status = MatchDelegate(fallback, by_prefix, file_location) - ->DeleteFile(file_location); - !status.has_value()) { + TaskGroup group; + group.SetExecutor(std::ref(delete_executor_)); + for (const auto& file_location : file_locations) { + group.Submit([file_io = MatchDelegate(fallback, by_prefix, file_location), + &file_location, &logger, &failed]() -> Status { + ScopedLogger bind(logger); + if (auto status = file_io->DeleteFile(file_location); !status.has_value()) { ICEBERG_LOG_WARN("Failed to delete {}: {}", file_location, status.error().message); ++failed; } - } - }; - // Plus the calling thread. Helpers log where the caller does. - auto logger = GetCurrentLogger(); - std::vector> helpers; - for (size_t i = 1; i < std::min(delete_threads_, file_locations.size()); ++i) { - helpers.push_back(std::async(std::launch::async, [&] { - ScopedLogger bind(logger); - delete_remaining(); - })); - } - delete_remaining(); - for (auto& helper : helpers) { - helper.get(); // Rethrows what a helper threw. + return {}; + }); } + ICEBERG_RETURN_UNEXPECTED(std::move(group).Run()); if (failed > 0) { return IOError("Failed to delete {} of {} files", failed.load(), file_locations.size()); @@ -396,15 +407,21 @@ Result> MakeS3FileIO( // Uses default credentials if properties are empty. ICEBERG_ASSIGN_OR_RAISE(auto fs, BuildArrowS3FileSystem(properties)); // Java defaults to one delete thread per processor. - size_t delete_threads = std::max(1u, std::thread::hardware_concurrency()); + int num_delete_threads = + std::max(1, static_cast(std::thread::hardware_concurrency())); if (const auto* value = FindProperty(properties, S3Properties::kDeleteNumThreads); value != nullptr) { - ICEBERG_ASSIGN_OR_RAISE(delete_threads, StringUtils::ParseNumber(*value)); - if (delete_threads == 0) { + ICEBERG_ASSIGN_OR_RAISE(num_delete_threads, StringUtils::ParseNumber(*value)); + if (num_delete_threads <= 0) { return InvalidArgument(R"("{}" must be positive)", S3Properties::kDeleteNumThreads); } } - return std::make_unique(std::move(fs), properties, delete_threads); + // The pool starts its workers on the first delete, so a FileIO that never + // deletes costs no threads. + ICEBERG_ARROW_ASSIGN_OR_RETURN(auto delete_pool, + ::arrow::internal::ThreadPool::Make(num_delete_threads)); + return std::make_unique(std::move(fs), properties, + std::move(delete_pool)); } Status FinalizeS3() { diff --git a/src/iceberg/arrow/s3/s3_properties.h b/src/iceberg/arrow/s3/s3_properties.h index bcc1865838..ad045f12fb 100644 --- a/src/iceberg/arrow/s3/s3_properties.h +++ b/src/iceberg/arrow/s3/s3_properties.h @@ -54,7 +54,8 @@ struct S3Properties { static constexpr std::string_view kConnectTimeoutMs = "s3.connect-timeout-ms"; /// Socket timeout in milliseconds static constexpr std::string_view kSocketTimeoutMs = "s3.socket-timeout-ms"; - /// Threads DeleteFiles uses; defaults to the hardware thread count, as in Java + /// Size of the pool DeleteFiles runs on; defaults to the hardware thread count, + /// as in Java static constexpr std::string_view kDeleteNumThreads = "s3.delete.num-threads"; }; From e12131404a47eb38a61be6dba8211c2130008523 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Sat, 10 Oct 2026 04:54:44 -0400 Subject: [PATCH 7/8] fix(io): configure shared executors for S3 deletes Use a shared caller-owned executor, defaulting to sequential deletion. Bound tasks, count delete exceptions, and drain accepted work on failure. --- mkdocs/docs/file-io.md | 21 +- src/iceberg/arrow/arrow_io_internal.h | 7 + src/iceberg/arrow/arrow_io_util.h | 9 + src/iceberg/arrow/s3/arrow_s3_file_io.cc | 160 +++++++++------- src/iceberg/arrow/s3/s3_properties.h | 3 - src/iceberg/test/arrow_s3_file_io_test.cc | 221 +++++++++++++++++++++- src/iceberg/util/executor.h | 4 +- 7 files changed, 342 insertions(+), 83 deletions(-) diff --git a/mkdocs/docs/file-io.md b/mkdocs/docs/file-io.md index 8064c8ae71..a914cb24aa 100644 --- a/mkdocs/docs/file-io.md +++ b/mkdocs/docs/file-io.md @@ -66,7 +66,6 @@ each file location's scheme. | `client.region` | `us-east-1` | Region to sign requests for | | `s3.endpoint` | `https://127.0.0.1:9000` | Endpoint to use instead of the AWS one. When absent, the `AWS_ENDPOINT_URL_S3` / `AWS_ENDPOINT_URL` environment variables are consulted | | `s3.path-style-access` | `true` | Address buckets as a path (`endpoint/bucket`) instead of a virtual host (`bucket.endpoint`). Only takes effect together with a custom endpoint | -| `s3.delete.num-threads` | `8` | Size of the thread pool `DeleteFiles` runs on, shared by every call on the FileIO. Defaults to the number of hardware threads | The following keys are specific to iceberg-cpp; they are not part of the Java Iceberg or REST specification property set: @@ -81,6 +80,26 @@ 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. 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 a6b85b6c98..401f107d0b 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 925ee12a27..98b57020e5 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/util/executor.h" namespace iceberg::arrow { @@ -47,6 +48,14 @@ 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. +/// 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 a0d18b776c..81f9da86a8 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -20,14 +20,15 @@ #include #include #include +#include #include +#include #include #include #include #include #include #include -#include #include #include #include @@ -35,7 +36,6 @@ #include #if ICEBERG_S3_ENABLED # include -# include #endif #include "iceberg/arrow/arrow_io_internal.h" @@ -47,10 +47,90 @@ #include "iceberg/util/macros.h" #include "iceberg/util/property_util.h" #include "iceberg/util/string_util.h" -#include "iceberg/util/task_group.h" namespace iceberg::arrow { +namespace { + +std::atomic delete_executor{nullptr}; + +struct DeleteTasks { + std::vector> futures; + + ~DeleteTasks() { + for (auto& future : futures) { + if (future.valid()) { + future.wait(); + } + } + } +}; + +} // 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 { + // Bound tasks and futures independently of the file count. + const auto workers = std::min(size_t{64}, file_locations.size()); + DeleteTasks tasks; + tasks.futures.reserve(workers); + for (size_t i = 0; i < workers; ++i) { + std::packaged_task task(work); + // Submit may accept the task before throwing. + tasks.futures.push_back(task.get_future()); + ExecutorTask executor_task([task = std::move(task)]() mutable { task(); }); + ICEBERG_RETURN_UNEXPECTED(executor->Submit(std::move(executor_task))); + } + for (auto& future : tasks.futures) { + future.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 { @@ -184,29 +264,12 @@ std::string CanonicalizeS3Scheme(std::string_view location) { return std::string(location); } -/// \brief Runs TaskGroup tasks on an Arrow thread pool. -class ThreadPoolExecutor final : public Executor { - public: - explicit ThreadPoolExecutor(std::shared_ptr<::arrow::internal::ThreadPool> pool) - : pool_(std::move(pool)) {} - - Status Submit(ExecutorTask task) override { - ICEBERG_ARROW_RETURN_NOT_OK(pool_->Spawn(std::move(task))); - return {}; - } - - private: - std::shared_ptr<::arrow::internal::ThreadPool> pool_; -}; - class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { public: ArrowS3FileIO(std::shared_ptr<::arrow::fs::FileSystem> arrow_fs, - std::unordered_map default_properties, - std::shared_ptr<::arrow::internal::ThreadPool> delete_pool) + std::unordered_map default_properties) : default_file_io_(std::make_shared(std::move(arrow_fs))), - default_properties_(std::move(default_properties)), - delete_executor_(std::move(delete_pool)) {} + default_properties_(std::move(default_properties)) {} Result> NewInputFile(std::string file_location) override; @@ -260,9 +323,6 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { std::shared_ptr default_file_io_; std::unordered_map default_properties_; - // Bounds the deletes of every DeleteFiles call on this FileIO, like the - // executor of Java's S3FileIO. - ThreadPoolExecutor delete_executor_; // Guards everything below; shared because reads happen per file operation. mutable std::shared_mutex mutex_; std::vector storage_credentials_; @@ -364,9 +424,7 @@ Status ArrowS3FileIO::DeleteFile(const std::string& file_location) { } Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations) { - // Like Java's S3FileIO: delete on the FileIO's bounded pool, keep going after - // a failure and report the count. Arrow has no batch delete for S3. One - // snapshot so the whole batch matches the same delegate generation. + // Keep one delegate snapshot for the whole batch. std::shared_ptr fallback; DelegatesByPrefix by_prefix; { @@ -374,30 +432,11 @@ Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations fallback = default_file_io_; by_prefix = file_io_by_prefix_; } - // Each failure is logged where the caller logs and only counted, so the - // returned error stays one line however large the batch. - auto logger = GetCurrentLogger(); - std::atomic failed = 0; - TaskGroup group; - group.SetExecutor(std::ref(delete_executor_)); - for (const auto& file_location : file_locations) { - group.Submit([file_io = MatchDelegate(fallback, by_prefix, file_location), - &file_location, &logger, &failed]() -> Status { - ScopedLogger bind(logger); - if (auto status = file_io->DeleteFile(file_location); !status.has_value()) { - ICEBERG_LOG_WARN("Failed to delete {}: {}", file_location, - status.error().message); - ++failed; - } - return {}; - }); - } - ICEBERG_RETURN_UNEXPECTED(std::move(group).Run()); - if (failed > 0) { - return IOError("Failed to delete {} of {} files", failed.load(), - file_locations.size()); - } - 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 @@ -406,22 +445,7 @@ Result> MakeS3FileIO( const std::unordered_map& properties) { // Uses default credentials if properties are empty. ICEBERG_ASSIGN_OR_RAISE(auto fs, BuildArrowS3FileSystem(properties)); - // Java defaults to one delete thread per processor. - int num_delete_threads = - std::max(1, static_cast(std::thread::hardware_concurrency())); - if (const auto* value = FindProperty(properties, S3Properties::kDeleteNumThreads); - value != nullptr) { - ICEBERG_ASSIGN_OR_RAISE(num_delete_threads, StringUtils::ParseNumber(*value)); - if (num_delete_threads <= 0) { - return InvalidArgument(R"("{}" must be positive)", S3Properties::kDeleteNumThreads); - } - } - // The pool starts its workers on the first delete, so a FileIO that never - // deletes costs no threads. - ICEBERG_ARROW_ASSIGN_OR_RETURN(auto delete_pool, - ::arrow::internal::ThreadPool::Make(num_delete_threads)); - return std::make_unique(std::move(fs), properties, - std::move(delete_pool)); + return std::make_unique(std::move(fs), properties); } Status FinalizeS3() { diff --git a/src/iceberg/arrow/s3/s3_properties.h b/src/iceberg/arrow/s3/s3_properties.h index ad045f12fb..50dafa56ff 100644 --- a/src/iceberg/arrow/s3/s3_properties.h +++ b/src/iceberg/arrow/s3/s3_properties.h @@ -54,9 +54,6 @@ struct S3Properties { static constexpr std::string_view kConnectTimeoutMs = "s3.connect-timeout-ms"; /// Socket timeout in milliseconds static constexpr std::string_view kSocketTimeoutMs = "s3.socket-timeout-ms"; - /// Size of the pool DeleteFiles runs on; defaults to the hardware thread count, - /// as in Java - static constexpr std::string_view kDeleteNumThreads = "s3.delete.num-threads"; }; /// \brief URI schemes served by the Arrow S3 FileIO, lower-case. diff --git a/src/iceberg/test/arrow_s3_file_io_test.cc b/src/iceberg/test/arrow_s3_file_io_test.cc index f75958c70c..7baffeb42f 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,193 @@ 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 run = std::async(std::launch::async, [&] { + 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); + EXPECT_THAT(status, HasErrorMessage("executor rejected deletion")); + EXPECT_EQ(attempted, 3); + EXPECT_EQ(executor.completed, failure == SubmitFailure::kAcceptedException ? 2 : 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,13 +491,24 @@ TEST_F(ArrowS3FileIOTest, RejectsIncompleteStaticCredentials) { "S3 client access key ID and secret access key must be set")); } -TEST_F(ArrowS3FileIOTest, RejectsInvalidDeleteThreads) { - for (std::string_view threads : {"0", "-1", "many"}) { - SCOPED_TRACE(threads); - EXPECT_THAT(MakeS3FileIO({{std::string(S3Properties::kDeleteNumThreads), - std::string(threads)}}), - IsError(ErrorKind::kInvalidArgument)); - } +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) { @@ -371,8 +574,8 @@ TEST_F(ArrowS3FileIOTest, DeleteFilesAttemptsEveryFile) { if (properties.empty()) { GTEST_SKIP() << "Set S3 properties to enable credential routing test"; } - // A thread per file, so the deletes run concurrently. - properties[std::string(S3Properties::kDeleteNumThreads)] = "3"; + test::ThreadExecutor executor; + ScopedDeleteExecutor configured(&executor); auto io_res = MakeS3FileIO(properties); ASSERT_THAT(io_res, IsOk()); diff --git a/src/iceberg/util/executor.h b/src/iceberg/util/executor.h index 749502b179..c16dc4b6cf 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,7 +41,7 @@ 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 +/// When an executor is configured, tasks may be called concurrently. Any /// shared mutable state captured by those callbacks must be synchronized by the caller. class ICEBERG_EXPORT Executor { public: From 4b042f839c4be5c2435acf0f50152554502ac80f Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Sat, 10 Oct 2026 03:02:49 -0700 Subject: [PATCH 8/8] fix(io): delete on the caller when the executor rejects a worker A rejected or throwing Submit returned its own error: with no worker accepted, no file was attempted, and otherwise the accepted workers still deleted every file while the failure count was lost. Log the rejection and let the calling thread take the files no accepted worker has, then report the count as usual. A rejected worker's future is still waited for, since Submit may accept a task before failing. Also name the worker bound, forward-declare Executor in arrow_io_util.h, and make the Executor docs say callbacks throughout. --- mkdocs/docs/file-io.md | 10 +++--- src/iceberg/arrow/arrow_io_util.h | 5 +-- src/iceberg/arrow/s3/arrow_s3_file_io.cc | 38 ++++++++++++++++++----- src/iceberg/test/arrow_s3_file_io_test.cc | 28 ++++++++++++++++- src/iceberg/util/executor.h | 4 +-- 5 files changed, 69 insertions(+), 16 deletions(-) diff --git a/mkdocs/docs/file-io.md b/mkdocs/docs/file-io.md index a914cb24aa..64b3682e7e 100644 --- a/mkdocs/docs/file-io.md +++ b/mkdocs/docs/file-io.md @@ -95,10 +95,12 @@ 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. 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. +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 diff --git a/src/iceberg/arrow/arrow_io_util.h b/src/iceberg/arrow/arrow_io_util.h index 98b57020e5..f5e4170fe9 100644 --- a/src/iceberg/arrow/arrow_io_util.h +++ b/src/iceberg/arrow/arrow_io_util.h @@ -29,7 +29,7 @@ #include "iceberg/file_io.h" #include "iceberg/iceberg_bundle_export.h" #include "iceberg/result.h" -#include "iceberg/util/executor.h" +#include "iceberg/type_fwd.h" namespace iceberg::arrow { @@ -52,7 +52,8 @@ ICEBERG_BUNDLE_EXPORT Result> MakeS3FileIO( /// /// 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 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); diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index 81f9da86a8..2028a254a2 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -44,6 +44,7 @@ #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" @@ -54,6 +55,9 @@ 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; @@ -66,6 +70,17 @@ struct DeleteTasks { } }; +// 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) { @@ -104,19 +119,28 @@ Status BulkDeleteFiles(const std::vector& file_locations, Executor* if (executor == nullptr) { work(); } else { - // Bound tasks and futures independently of the file count. - const auto workers = std::min(size_t{64}, file_locations.size()); + const auto workers = std::min(kMaxDeleteWorkers, file_locations.size()); DeleteTasks tasks; tasks.futures.reserve(workers); - for (size_t i = 0; i < workers; ++i) { + size_t accepted = 0; + for (; accepted < workers; ++accepted) { std::packaged_task task(work); - // Submit may accept the task before throwing. + // Submit may accept the task before failing. tasks.futures.push_back(task.get_future()); ExecutorTask executor_task([task = std::move(task)]() mutable { task(); }); - ICEBERG_RETURN_UNEXPECTED(executor->Submit(std::move(executor_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; + } } - for (auto& future : tasks.futures) { - future.get(); + // 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) { diff --git a/src/iceberg/test/arrow_s3_file_io_test.cc b/src/iceberg/test/arrow_s3_file_io_test.cc index 7baffeb42f..19d66fe970 100644 --- a/src/iceberg/test/arrow_s3_file_io_test.cc +++ b/src/iceberg/test/arrow_s3_file_io_test.cc @@ -290,7 +290,9 @@ TEST(BulkDeleteFilesTest, DrainsAcceptedTasksWhenSubmissionFails) { 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 {}; @@ -310,12 +312,36 @@ TEST(BulkDeleteFilesTest, DrainsAcceptedTasksWhenSubmissionFails) { EXPECT_EQ(submitted, std::future_status::ready); EXPECT_EQ(first, std::future_status::ready); EXPECT_EQ(returned, std::future_status::timeout); - EXPECT_THAT(status, HasErrorMessage("executor rejected deletion")); + // 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}) { diff --git a/src/iceberg/util/executor.h b/src/iceberg/util/executor.h index c16dc4b6cf..625c833f12 100644 --- a/src/iceberg/util/executor.h +++ b/src/iceberg/util/executor.h @@ -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, tasks 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;