|
| 1 | +/* |
| 2 | + * Licensed to the Apache Software Foundation (ASF) under one |
| 3 | + * or more contributor license agreements. See the NOTICE file |
| 4 | + * distributed with this work for additional information |
| 5 | + * regarding copyright ownership. The ASF licenses this file |
| 6 | + * to you under the Apache License, Version 2.0 (the |
| 7 | + * "License"); you may not use this file except in compliance |
| 8 | + * with the License. You may obtain a copy of the License at |
| 9 | + * |
| 10 | + * http://www.apache.org/licenses/LICENSE-2.0 |
| 11 | + * |
| 12 | + * Unless required by applicable law or agreed to in writing, |
| 13 | + * software distributed under the License is distributed on an |
| 14 | + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| 15 | + * KIND, either express or implied. See the License for the |
| 16 | + * specific language governing permissions and limitations |
| 17 | + * under the License. |
| 18 | + */ |
| 19 | + |
| 20 | +#include "iceberg/data/position_delete_update.h" |
| 21 | + |
| 22 | +#include <algorithm> |
| 23 | +#include <format> |
| 24 | +#include <iterator> |
| 25 | +#include <map> |
| 26 | +#include <memory> |
| 27 | +#include <optional> |
| 28 | +#include <set> |
| 29 | +#include <string> |
| 30 | +#include <string_view> |
| 31 | +#include <tuple> |
| 32 | +#include <unordered_map> |
| 33 | +#include <utility> |
| 34 | +#include <vector> |
| 35 | + |
| 36 | +#include "iceberg/data/delete_loader.h" |
| 37 | +#include "iceberg/data/position_delete_writer.h" |
| 38 | +#include "iceberg/deletes/dv_writer.h" |
| 39 | +#include "iceberg/file_format.h" |
| 40 | +#include "iceberg/file_io.h" |
| 41 | +#include "iceberg/location_provider.h" |
| 42 | +#include "iceberg/manifest/manifest_entry.h" |
| 43 | +#include "iceberg/partition_spec.h" |
| 44 | +#include "iceberg/schema.h" |
| 45 | +#include "iceberg/table.h" |
| 46 | +#include "iceberg/table_metadata.h" |
| 47 | +#include "iceberg/table_properties.h" |
| 48 | +#include "iceberg/table_scan.h" |
| 49 | +#include "iceberg/update/row_delta.h" |
| 50 | +#include "iceberg/util/macros.h" |
| 51 | +#include "iceberg/util/string_util.h" |
| 52 | +#include "iceberg/util/uuid.h" |
| 53 | + |
| 54 | +namespace iceberg { |
| 55 | + |
| 56 | +namespace { |
| 57 | + |
| 58 | +struct TargetFile { |
| 59 | + std::shared_ptr<DataFile> data_file; |
| 60 | + std::shared_ptr<PartitionSpec> spec; |
| 61 | + std::vector<std::shared_ptr<DataFile>> position_delete_files; |
| 62 | +}; |
| 63 | + |
| 64 | +struct PreparedDeletes { |
| 65 | + std::vector<std::shared_ptr<DataFile>> added_files; |
| 66 | + std::vector<std::shared_ptr<DataFile>> rewritten_files; |
| 67 | + std::vector<std::string> referenced_data_files; |
| 68 | +}; |
| 69 | + |
| 70 | +} // namespace |
| 71 | + |
| 72 | +class PositionDeleteUpdate::Impl { |
| 73 | + public: |
| 74 | + explicit Impl(std::shared_ptr<Table> table) : table_(std::move(table)) {} |
| 75 | + |
| 76 | + Status Delete(std::string_view data_file_path, int64_t pos) { |
| 77 | + ICEBERG_PRECHECK(!terminal_, "Position delete update is no longer usable"); |
| 78 | + ICEBERG_PRECHECK(!data_file_path.empty(), "Data file path cannot be empty"); |
| 79 | + ICEBERG_PRECHECK(pos >= 0, "Position delete must be non-negative: {}", pos); |
| 80 | + deletes_[std::string(data_file_path)].push_back(pos); |
| 81 | + return {}; |
| 82 | + } |
| 83 | + |
| 84 | + Status Commit() { |
| 85 | + ICEBERG_PRECHECK(!terminal_, "Position delete update is no longer usable"); |
| 86 | + ICEBERG_PRECHECK(!deletes_.empty(), "Position delete update is empty"); |
| 87 | + ICEBERG_PRECHECK(table_->metadata()->format_version >= 2, |
| 88 | + "Position deletes require table format version 2 or later"); |
| 89 | + ICEBERG_RETURN_UNEXPECTED(CleanupOutput()); |
| 90 | + |
| 91 | + ICEBERG_ASSIGN_OR_RAISE(auto snapshot, table_->current_snapshot()); |
| 92 | + ICEBERG_ASSIGN_OR_RAISE(auto targets, ResolveTargets()); |
| 93 | + |
| 94 | + auto prepared = table_->metadata()->format_version >= 3 |
| 95 | + ? WriteDeletionVectors(targets) |
| 96 | + : WriteParquetDeletes(targets); |
| 97 | + if (!prepared.has_value()) { |
| 98 | + return FailAfterCleanup(std::move(prepared.error())); |
| 99 | + } |
| 100 | + |
| 101 | + auto row_delta_result = table_->NewRowDelta(); |
| 102 | + if (!row_delta_result.has_value()) { |
| 103 | + return FailAfterCleanup(std::move(row_delta_result.error())); |
| 104 | + } |
| 105 | + auto row_delta = std::move(row_delta_result.value()); |
| 106 | + row_delta->ValidateFromSnapshot(snapshot->snapshot_id) |
| 107 | + .ValidateDataFilesExist(prepared->referenced_data_files) |
| 108 | + .ValidateDeletedFiles(); |
| 109 | + for (const auto& file : prepared->added_files) { |
| 110 | + row_delta->AddDeletes(file); |
| 111 | + } |
| 112 | + for (const auto& file : prepared->rewritten_files) { |
| 113 | + row_delta->RemoveDeletes(file); |
| 114 | + } |
| 115 | + |
| 116 | + auto status = row_delta->Commit(); |
| 117 | + if (!status.has_value()) { |
| 118 | + if (status.error().kind == ErrorKind::kCommitStateUnknown) { |
| 119 | + terminal_ = true; |
| 120 | + output_paths_.clear(); |
| 121 | + } else { |
| 122 | + return FailAfterCleanup(std::move(status.error())); |
| 123 | + } |
| 124 | + return status; |
| 125 | + } |
| 126 | + |
| 127 | + terminal_ = true; |
| 128 | + output_paths_.clear(); |
| 129 | + return {}; |
| 130 | + } |
| 131 | + |
| 132 | + private: |
| 133 | + Result<std::unordered_map<std::string, TargetFile>> ResolveTargets() const { |
| 134 | + ICEBERG_ASSIGN_OR_RAISE(auto scan_builder, table_->NewScan()); |
| 135 | + ICEBERG_ASSIGN_OR_RAISE(auto scan, scan_builder->Build()); |
| 136 | + ICEBERG_ASSIGN_OR_RAISE(auto tasks, scan->PlanFiles()); |
| 137 | + |
| 138 | + std::unordered_map<std::string, TargetFile> targets; |
| 139 | + for (const auto& task : tasks) { |
| 140 | + const auto& data_file = task->data_file(); |
| 141 | + if (!deletes_.contains(data_file->file_path)) { |
| 142 | + continue; |
| 143 | + } |
| 144 | + |
| 145 | + ICEBERG_PRECHECK(data_file->partition_spec_id.has_value(), |
| 146 | + "Data file is missing partition spec ID: {}", |
| 147 | + data_file->file_path); |
| 148 | + ICEBERG_ASSIGN_OR_RAISE(auto spec, table_->metadata()->PartitionSpecById( |
| 149 | + *data_file->partition_spec_id)); |
| 150 | + |
| 151 | + TargetFile target{.data_file = data_file, .spec = std::move(spec)}; |
| 152 | + for (const auto& delete_file : task->delete_files()) { |
| 153 | + if (delete_file->content == DataFile::Content::kPositionDeletes) { |
| 154 | + target.position_delete_files.push_back(delete_file); |
| 155 | + } |
| 156 | + } |
| 157 | + |
| 158 | + auto [_, inserted] = targets.emplace(data_file->file_path, std::move(target)); |
| 159 | + ICEBERG_PRECHECK(inserted, "Duplicate live data file path: {}", |
| 160 | + data_file->file_path); |
| 161 | + } |
| 162 | + |
| 163 | + for (const auto& [path, _] : deletes_) { |
| 164 | + ICEBERG_PRECHECK(targets.contains(path), "Cannot find live data file: {}", path); |
| 165 | + } |
| 166 | + return targets; |
| 167 | + } |
| 168 | + |
| 169 | + Result<PreparedDeletes> WriteDeletionVectors( |
| 170 | + const std::unordered_map<std::string, TargetFile>& targets) { |
| 171 | + ICEBERG_ASSIGN_OR_RAISE(auto location_provider, table_->location_provider()); |
| 172 | + auto output_path = location_provider->NewDataLocation( |
| 173 | + std::format("position-deletes-{}.puffin", Uuid::GenerateV7().ToString())); |
| 174 | + output_paths_.insert(output_path); |
| 175 | + |
| 176 | + DeleteLoader loader(table_->io()); |
| 177 | + ICEBERG_ASSIGN_OR_RAISE( |
| 178 | + auto writer, |
| 179 | + DVWriter::Make(DVWriterOptions{ |
| 180 | + .path = output_path, |
| 181 | + .io = table_->io(), |
| 182 | + .load_previous_deletes = [&targets, &loader](std::string_view path) |
| 183 | + -> Result<std::optional<PositionDeleteIndex>> { |
| 184 | + const auto target = targets.find(std::string(path)); |
| 185 | + ICEBERG_CHECK(target != targets.end(), |
| 186 | + "Missing target data file while loading deletes: {}", path); |
| 187 | + if (target->second.position_delete_files.empty()) { |
| 188 | + return std::nullopt; |
| 189 | + } |
| 190 | + ICEBERG_ASSIGN_OR_RAISE( |
| 191 | + auto index, |
| 192 | + loader.LoadPositionDeletes(target->second.position_delete_files, path)); |
| 193 | + return std::optional<PositionDeleteIndex>(std::move(index)); |
| 194 | + }, |
| 195 | + })); |
| 196 | + |
| 197 | + for (const auto& [path, positions] : deletes_) { |
| 198 | + const auto& target = targets.at(path); |
| 199 | + for (int64_t pos : positions) { |
| 200 | + ICEBERG_RETURN_UNEXPECTED( |
| 201 | + writer->Delete(path, pos, target.spec, target.data_file->partition)); |
| 202 | + } |
| 203 | + } |
| 204 | + ICEBERG_RETURN_UNEXPECTED(writer->Close()); |
| 205 | + ICEBERG_ASSIGN_OR_RAISE(auto result, writer->Metadata()); |
| 206 | + return PreparedDeletes{ |
| 207 | + .added_files = std::move(result.data_files), |
| 208 | + .rewritten_files = std::move(result.rewritten_delete_files), |
| 209 | + .referenced_data_files = std::move(result.referenced_data_files), |
| 210 | + }; |
| 211 | + } |
| 212 | + |
| 213 | + Result<PreparedDeletes> WriteParquetDeletes( |
| 214 | + const std::unordered_map<std::string, TargetFile>& targets) { |
| 215 | + ICEBERG_ASSIGN_OR_RAISE(auto location_provider, table_->location_provider()); |
| 216 | + ICEBERG_ASSIGN_OR_RAISE(auto schema, table_->schema()); |
| 217 | + |
| 218 | + PreparedDeletes result; |
| 219 | + result.added_files.reserve(deletes_.size()); |
| 220 | + result.referenced_data_files.reserve(deletes_.size()); |
| 221 | + auto properties = table_->properties().configs(); |
| 222 | + properties[TableProperties::kParquetCompression.key()] = |
| 223 | + table_->properties().Get(TableProperties::kDeleteParquetCompression); |
| 224 | + properties[TableProperties::kParquetCompressionLevel.key()] = |
| 225 | + table_->properties().Get(TableProperties::kDeleteParquetCompressionLevel); |
| 226 | + const auto write_uuid = Uuid::GenerateV7().ToString(); |
| 227 | + size_t file_number = 0; |
| 228 | + |
| 229 | + for (const auto& [path, positions] : deletes_) { |
| 230 | + auto sorted_positions = positions; |
| 231 | + std::ranges::sort(sorted_positions); |
| 232 | + |
| 233 | + const auto& target = targets.at(path); |
| 234 | + const auto filename = |
| 235 | + std::format("position-deletes-{}-{}.parquet", write_uuid, file_number++); |
| 236 | + auto output_path = location_provider->NewDataLocation(filename); |
| 237 | + output_paths_.insert(output_path); |
| 238 | + |
| 239 | + ICEBERG_ASSIGN_OR_RAISE(auto writer, |
| 240 | + PositionDeleteWriter::Make(PositionDeleteWriterOptions{ |
| 241 | + .path = output_path, |
| 242 | + .schema = schema, |
| 243 | + .spec = target.spec, |
| 244 | + .partition = target.data_file->partition, |
| 245 | + .format = FileFormatType::kParquet, |
| 246 | + .io = table_->io(), |
| 247 | + .properties = properties, |
| 248 | + })); |
| 249 | + for (int64_t pos : sorted_positions) { |
| 250 | + ICEBERG_RETURN_UNEXPECTED(writer->WriteDelete(path, pos)); |
| 251 | + } |
| 252 | + ICEBERG_RETURN_UNEXPECTED(writer->Close()); |
| 253 | + ICEBERG_ASSIGN_OR_RAISE(auto metadata, writer->Metadata()); |
| 254 | + result.added_files.insert(result.added_files.end(), |
| 255 | + std::make_move_iterator(metadata.data_files.begin()), |
| 256 | + std::make_move_iterator(metadata.data_files.end())); |
| 257 | + result.referenced_data_files.push_back(path); |
| 258 | + } |
| 259 | + return result; |
| 260 | + } |
| 261 | + |
| 262 | + Status CleanupOutput() { |
| 263 | + std::optional<Error> first_error; |
| 264 | + for (auto it = output_paths_.begin(); it != output_paths_.end();) { |
| 265 | + auto status = table_->io()->DeleteFile(*it); |
| 266 | + if (status.has_value()) { |
| 267 | + it = output_paths_.erase(it); |
| 268 | + continue; |
| 269 | + } |
| 270 | + |
| 271 | + if (!first_error.has_value()) { |
| 272 | + first_error = std::move(status.error()); |
| 273 | + } else { |
| 274 | + first_error->message += "; additionally failed to delete output file: "; |
| 275 | + first_error->message += status.error().message; |
| 276 | + } |
| 277 | + ++it; |
| 278 | + } |
| 279 | + if (first_error.has_value()) { |
| 280 | + return std::unexpected(std::move(*first_error)); |
| 281 | + } |
| 282 | + return {}; |
| 283 | + } |
| 284 | + |
| 285 | + Status FailAfterCleanup(Error error) { |
| 286 | + auto cleanup_status = CleanupOutput(); |
| 287 | + if (!cleanup_status.has_value()) { |
| 288 | + error.message += "; additionally failed to clean output files: "; |
| 289 | + error.message += cleanup_status.error().message; |
| 290 | + } |
| 291 | + return std::unexpected(std::move(error)); |
| 292 | + } |
| 293 | + |
| 294 | + std::shared_ptr<Table> table_; |
| 295 | + std::map<std::string, std::vector<int64_t>, StringLess> deletes_; |
| 296 | + std::set<std::string, StringLess> output_paths_; |
| 297 | + bool terminal_ = false; |
| 298 | +}; |
| 299 | + |
| 300 | +PositionDeleteUpdate::PositionDeleteUpdate(std::unique_ptr<Impl> impl) |
| 301 | + : impl_(std::move(impl)) {} |
| 302 | + |
| 303 | +PositionDeleteUpdate::~PositionDeleteUpdate() = default; |
| 304 | + |
| 305 | +Result<std::unique_ptr<PositionDeleteUpdate>> PositionDeleteUpdate::Make( |
| 306 | + std::shared_ptr<Table> table) { |
| 307 | + ICEBERG_PRECHECK(table != nullptr, |
| 308 | + "Cannot create position delete update without table"); |
| 309 | + return std::unique_ptr<PositionDeleteUpdate>( |
| 310 | + new PositionDeleteUpdate(std::make_unique<Impl>(std::move(table)))); |
| 311 | +} |
| 312 | + |
| 313 | +PositionDeleteUpdate& PositionDeleteUpdate::Delete(std::string_view data_file_path, |
| 314 | + int64_t pos) { |
| 315 | + ICEBERG_BUILDER_RETURN_IF_ERROR(impl_->Delete(data_file_path, pos)); |
| 316 | + return *this; |
| 317 | +} |
| 318 | + |
| 319 | +Status PositionDeleteUpdate::Commit() { |
| 320 | + ICEBERG_RETURN_UNEXPECTED(CheckErrors()); |
| 321 | + return impl_->Commit(); |
| 322 | +} |
| 323 | + |
| 324 | +} // namespace iceberg |
0 commit comments