Delete multiple objects at once
| 2736 | |
| 2737 | // Delete multiple objects at once |
| 2738 | Future<> DeleteObjectsAsync(const std::string& bucket, |
| 2739 | const std::vector<std::string>& keys) { |
| 2740 | struct DeleteCallback { |
| 2741 | std::string bucket; |
| 2742 | |
| 2743 | Status operator()(const S3Model::DeleteObjectsOutcome& outcome) const { |
| 2744 | if (!outcome.IsSuccess()) { |
| 2745 | return ErrorToStatus("DeleteObjects", outcome.GetError()); |
| 2746 | } |
| 2747 | // Also need to check per-key errors, even on successful outcome |
| 2748 | // See |
| 2749 | // https://docs.aws.amazon.com/fr_fr/AmazonS3/latest/API/multiobjectdeleteapi.html |
| 2750 | const auto& errors = outcome.GetResult().GetErrors(); |
| 2751 | if (!errors.empty()) { |
| 2752 | std::stringstream ss; |
| 2753 | ss << "Got the following " << errors.size() |
| 2754 | << " errors when deleting objects in S3 bucket '" << bucket << "':\n"; |
| 2755 | for (const auto& error : errors) { |
| 2756 | ss << "- key '" << error.GetKey() << "': " << error.GetMessage() << "\n"; |
| 2757 | } |
| 2758 | return Status::IOError(ss.str()); |
| 2759 | } |
| 2760 | return Status::OK(); |
| 2761 | } |
| 2762 | }; |
| 2763 | |
| 2764 | const auto chunk_size = static_cast<size_t>(kMultipleDeleteMaxKeys); |
| 2765 | const DeleteCallback delete_cb{bucket}; |
| 2766 | |
| 2767 | std::vector<Future<>> futures; |
| 2768 | futures.reserve(bit_util::CeilDiv(keys.size(), chunk_size)); |
| 2769 | |
| 2770 | for (size_t start = 0; start < keys.size(); start += chunk_size) { |
| 2771 | S3Model::DeleteObjectsRequest req; |
| 2772 | S3Model::Delete del; |
| 2773 | size_t remaining = keys.size() - start; |
| 2774 | size_t next_chunk_size = std::min(remaining, chunk_size); |
| 2775 | for (size_t i = start; i < start + next_chunk_size; ++i) { |
| 2776 | del.AddObjects(S3Model::ObjectIdentifier().WithKey(ToAwsString(keys[i]))); |
| 2777 | } |
| 2778 | req.SetBucket(ToAwsString(bucket)); |
| 2779 | req.SetDelete(std::move(del)); |
| 2780 | ARROW_ASSIGN_OR_RAISE( |
| 2781 | auto fut, |
| 2782 | SubmitIO(io_context_, |
| 2783 | [holder = holder_, req = std::move(req), delete_cb]() -> Status { |
| 2784 | ARROW_ASSIGN_OR_RAISE(auto client_lock, holder->Lock()); |
| 2785 | return delete_cb(client_lock.Move()->DeleteObjects(req)); |
| 2786 | })); |
| 2787 | futures.push_back(std::move(fut)); |
| 2788 | } |
| 2789 | |
| 2790 | return AllFinished(futures); |
| 2791 | } |
| 2792 | |
| 2793 | Status DeleteObjects(const std::string& bucket, const std::vector<std::string>& keys) { |
| 2794 | return DeleteObjectsAsync(bucket, keys).status(); |
no test coverage detected