Delete multiple objects at once
| 2795 | |
| 2796 | // Delete multiple objects at once |
| 2797 | Future<> DeleteObjectsAsync(const std::string& bucket, |
| 2798 | const std::vector<std::string>& keys) { |
| 2799 | struct DeleteCallback { |
| 2800 | std::string bucket; |
| 2801 | |
| 2802 | Status operator()(const S3Model::DeleteObjectsOutcome& outcome) const { |
| 2803 | if (!outcome.IsSuccess()) { |
| 2804 | return ErrorToStatus("DeleteObjects", outcome.GetError()); |
| 2805 | } |
| 2806 | // Also need to check per-key errors, even on successful outcome |
| 2807 | // See |
| 2808 | // https://docs.aws.amazon.com/fr_fr/AmazonS3/latest/API/multiobjectdeleteapi.html |
| 2809 | const auto& errors = outcome.GetResult().GetErrors(); |
| 2810 | if (!errors.empty()) { |
| 2811 | std::stringstream ss; |
| 2812 | ss << "Got the following " << errors.size() |
| 2813 | << " errors when deleting objects in S3 bucket '" << bucket << "':\n"; |
| 2814 | for (const auto& error : errors) { |
| 2815 | ss << "- key '" << error.GetKey() << "': " << error.GetMessage() << "\n"; |
| 2816 | } |
| 2817 | return Status::IOError(ss.str()); |
| 2818 | } |
| 2819 | return Status::OK(); |
| 2820 | } |
| 2821 | }; |
| 2822 | |
| 2823 | const auto chunk_size = static_cast<size_t>(kMultipleDeleteMaxKeys); |
| 2824 | const DeleteCallback delete_cb{bucket}; |
| 2825 | |
| 2826 | std::vector<Future<>> futures; |
| 2827 | futures.reserve(bit_util::CeilDiv(keys.size(), chunk_size)); |
| 2828 | |
| 2829 | for (size_t start = 0; start < keys.size(); start += chunk_size) { |
| 2830 | S3Model::DeleteObjectsRequest req; |
| 2831 | S3Model::Delete del; |
| 2832 | size_t remaining = keys.size() - start; |
| 2833 | size_t next_chunk_size = std::min(remaining, chunk_size); |
| 2834 | for (size_t i = start; i < start + next_chunk_size; ++i) { |
| 2835 | del.AddObjects(S3Model::ObjectIdentifier().WithKey(ToAwsString(keys[i]))); |
| 2836 | } |
| 2837 | req.SetBucket(ToAwsString(bucket)); |
| 2838 | req.SetDelete(std::move(del)); |
| 2839 | ARROW_ASSIGN_OR_RAISE( |
| 2840 | auto fut, |
| 2841 | SubmitIO(io_context_, |
| 2842 | [holder = holder_, req = std::move(req), delete_cb]() -> Status { |
| 2843 | ARROW_ASSIGN_OR_RAISE(auto client_lock, holder->Lock()); |
| 2844 | return delete_cb(client_lock.Move()->DeleteObjects(req)); |
| 2845 | })); |
| 2846 | futures.push_back(std::move(fut)); |
| 2847 | } |
| 2848 | |
| 2849 | return AllFinished(futures); |
| 2850 | } |
| 2851 | |
| 2852 | Status DeleteObjects(const std::string& bucket, const std::vector<std::string>& keys) { |
| 2853 | return DeleteObjectsAsync(bucket, keys).status(); |
no test coverage detected