Repository navigation
fix(io): attempt every S3 delete and report how many failed - #966
plusplusjiajia wants to merge 6 commits into
Conversation
fee511b to
1f5fb83
Compare
| size_t failed = 0; | ||
| for (const auto& file_location : file_locations) { | ||
| locations_by_io[&FileIOForPath(file_location)].push_back(file_location); | ||
| if (auto status = FileIOForPath(file_location).DeleteFile(file_location); |
There was a problem hiding this comment.
Should we use a thread pool to delete files concurrently? Sequential S3 deletes could be slow for large batches. Java's S3FileIO does that, so maybe we should pursue that too, but not a requirement for this PR.
There was a problem hiding this comment.
@zhjwpku Thanks! Done: deletes now run on up to s3.delete.num-threads threads (default: hardware threads), as in Java. Java also batches keys into DeleteObjects; Arrow has no public batch delete for S3, so each thread deletes one file at a time.
| auto logger = GetCurrentLogger(); | ||
| std::vector<std::future<void>> helpers; | ||
| for (size_t i = 1; i < std::min(delete_threads_, file_locations.size()); ++i) { | ||
| helpers.push_back(std::async(std::launch::async, [&] { |
There was a problem hiding this comment.
We have TaskGroup and Executor abstractions, not sure if they fit it here.
There was a problem hiding this comment.
Maybe not since there is no need to retry here. cc @HuaHuaY
There was a problem hiding this comment.
@zhjwpku Thanks! TaskGroup needs a caller-supplied Executor, which DeleteFiles can't receive, and ExpireSnapshots already retries it. So plain threads here.
There was a problem hiding this comment.
@plusplusjiajia
I get your opinion above.
In this case, could ArrowS3FileIO has its own bounded executor initialized from s3.delete.num-threads and reuse it across DeleteFiles calls?
This enforces the concurrency limit across calls and allow using TaskGroup instead of creating threads per invocation.
There was a problem hiding this comment.
@kamcheungting-db Thanks! The FileIO now owns a bounded Arrow thread pool shared across DeleteFiles calls, with tasks run through TaskGroup.
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.
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.
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.
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.
41ff77e to
492c6cb
Compare
| : default_file_io_(std::make_shared<ArrowFileSystemFileIO>(std::move(arrow_fs))), | ||
| default_properties_(std::move(default_properties)) {} | ||
| default_properties_(std::move(default_properties)), | ||
| delete_threads_(delete_threads) {} |
There was a problem hiding this comment.
| delete_threads_(delete_threads) {} | |
| num_delete_threads_(delete_threads) {} |
There was a problem hiding this comment.
Thanks! That member was removed; the FileIO now owns the thread pool.
|
|
||
| std::shared_ptr<ArrowFileSystemFileIO> default_file_io_; | ||
| std::unordered_map<std::string, std::string> default_properties_; | ||
| size_t delete_threads_; |
There was a problem hiding this comment.
| size_t delete_threads_; | |
| size_t num_delete_threads_; |
There was a problem hiding this comment.
Thanks! That member was removed; the FileIO now owns the thread pool.
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.
| } | ||
| // 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, |
There was a problem hiding this comment.
I don't think each FileIO should create its own thread pool. Java's executor is static and shared across instances. Could we instead expose a static/public configuration API for a caller-owned iceberg::Executor*, defaulting to nullptr for sequential deletion? This should also work for registry-created FileIOs. Please document initialization and executor lifetime requirements; the executor, rather than each FileIO, should control concurrency.
| } | ||
| group.Submit([file_io = MatchDelegate(fallback, by_prefix, file_location), | ||
| &file_location, &logger, &failed]() -> Status { | ||
| ScopedLogger bind(logger); |
There was a problem hiding this comment.
Please catch exceptions inside each delete task and include them in the failure count. TaskGroup currently invokes the task without catching exceptions, and Arrow's worker loop does not catch them either, so a throwing delete can terminate the process instead of returning a Status. Please cover this with a throwing-task test and verify that the remaining files are still attempted.
| } | ||
| for (auto& [file_io, locations] : locations_by_io) { | ||
| ICEBERG_RETURN_UNEXPECTED(file_io->DeleteFiles(locations)); | ||
| ICEBERG_RETURN_UNEXPECTED(std::move(group).Run()); |
There was a problem hiding this comment.
Please ensure every accepted task completes before DeleteFiles exits, including when submission throws after earlier tasks were accepted. TaskGroup uses promise-backed futures whose destruction does not wait, so unwinding can leave queued tasks referencing destroyed logger, failed, or input data. Please add exception-safe draining and a test where submission fails after some work has been accepted.
| } else { | ||
| it->second.push_back(file_location); | ||
| } | ||
| group.Submit([file_io = MatchDelegate(fallback, by_prefix, file_location), |
There was a problem hiding this comment.
Could we bound the number of in-flight tasks rather than creating a task and future for every file? Large snapshot expirations can contain millions of paths, and TaskGroup retains all futures until completion. A bounded submission window or chunked tasks would avoid this allocation and queue growth while preserving the delegate snapshot and per-file failure accounting.
| | `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 | |
There was a problem hiding this comment.
As my other comment, this should not be a S3FileIO property.
Follow-up to #898 (comment).
ArrowS3FileIO::DeleteFilesreturned at the first credential prefix whose delete failed, so files under the remaining prefixes were never attempted. Java'sS3FileIO.deleteFilesattempts every batch on a thread pool, logs each failed path, and throwsBulkDeletionFailureExceptionwith the count.This matches that:
Failed to delete N of M files.s3.delete.num-threads(default: the number of hardware threads) and shared by every call, as Java'sS3FileIOexecutor is. The pool starts its workers on the first delete.Java also packs keys into
DeleteObjectsbatches. Arrow has no public batch delete for S3, so each thread deletes one file at a time, as Arrow's ownDeleteFilesdoes.ResolvingFileIOstays fail-fast across delegates, as Java's does.DeleteFilesAttemptsEveryFileputs an allowed file between two denied ones; it fails onmainand passes here.