From 2e7240e6204a923b68dd7af94bc98b6dd5575d24 Mon Sep 17 00:00:00 2001 From: xuanyili Date: Sat, 5 Sep 2026 21:44:05 +0000 Subject: [PATCH] feat: add position_deletes metadata table Reuse one reader and projection context per scan. Append and decode decimal partitions as Decimal128, and encode Avro fixed decimals using their schema width rather than all 16 Arrow bytes. Preserve logical date storage values and cover the complete partition write/read path. --- src/iceberg/CMakeLists.txt | 1 + src/iceberg/arrow_c_data_util.cc | 6 + src/iceberg/arrow_c_data_util_internal.h | 8 + src/iceberg/arrow_row_builder.cc | 10 + src/iceberg/arrow_row_builder_internal.h | 4 + src/iceberg/avro/avro_data_util.cc | 1 + src/iceberg/avro/avro_direct_encoder.cc | 4 +- src/iceberg/deletes/position_delete_index.cc | 4 + src/iceberg/deletes/position_delete_index.h | 4 + src/iceberg/file_reader.h | 6 + src/iceberg/inspect/metadata_table.h | 1 + src/iceberg/inspect/position_deletes_table.cc | 767 ++++++++++++++++++ src/iceberg/inspect/position_deletes_table.h | 72 ++ src/iceberg/manifest/manifest_adapter.cc | 5 +- src/iceberg/manifest/manifest_reader.cc | 12 + src/iceberg/parquet/parquet_reader.cc | 13 + src/iceberg/parquet/parquet_reader.h | 2 + src/iceberg/test/CMakeLists.txt | 3 + src/iceberg/test/avro_data_test.cc | 7 +- .../test/position_deletes_table_test.cc | 701 ++++++++++++++++ src/iceberg/type_fwd.h | 1 + 21 files changed, 1625 insertions(+), 7 deletions(-) create mode 100644 src/iceberg/inspect/position_deletes_table.cc create mode 100644 src/iceberg/inspect/position_deletes_table.h create mode 100644 src/iceberg/test/position_deletes_table_test.cc diff --git a/src/iceberg/CMakeLists.txt b/src/iceberg/CMakeLists.txt index 8a98274ff..2457f8eed 100644 --- a/src/iceberg/CMakeLists.txt +++ b/src/iceberg/CMakeLists.txt @@ -55,6 +55,7 @@ set(ICEBERG_SOURCES geospatial.cc inspect/history_table.cc inspect/metadata_table.cc + inspect/position_deletes_table.cc inspect/snapshots_table.cc inheritable_metadata.cc json_serde.cc diff --git a/src/iceberg/arrow_c_data_util.cc b/src/iceberg/arrow_c_data_util.cc index a77012a69..8139ad0da 100644 --- a/src/iceberg/arrow_c_data_util.cc +++ b/src/iceberg/arrow_c_data_util.cc @@ -220,6 +220,12 @@ Status AppendValue(const ArrowSchema& input_schema, const ArrowArray& input_arra } // namespace +Status AppendArrayValue(const ArrowSchema& input_schema, const ArrowArray& input_array, + const ArrowArrayView& input_view, int64_t row_index, + ArrowArray* output_array) { + return AppendValue(input_schema, input_array, input_view, row_index, output_array); +} + ProjectionContext::ProjectionContext(ProjectionContext&& other) noexcept : input_schema_(std::exchange(other.input_schema_, nullptr)), output_schema_(std::exchange(other.output_schema_, nullptr)), diff --git a/src/iceberg/arrow_c_data_util_internal.h b/src/iceberg/arrow_c_data_util_internal.h index e02db29a6..709d72ea0 100644 --- a/src/iceberg/arrow_c_data_util_internal.h +++ b/src/iceberg/arrow_c_data_util_internal.h @@ -35,6 +35,8 @@ #include "iceberg/result.h" #include "iceberg/type_fwd.h" +struct ArrowArrayView; + namespace iceberg { /// \brief Cached state for ProjectBatch over one input/output schema pair. @@ -242,4 +244,10 @@ ICEBERG_EXPORT Result ProjectBatch(ArrowArray* input_batch, std::span row_indices, ProjectionContext& projection); +/// \brief Append one value from an Arrow array into a compatible nanoarrow builder. +ICEBERG_EXPORT Status AppendArrayValue(const ArrowSchema& input_schema, + const ArrowArray& input_array, + const ArrowArrayView& input_view, + int64_t row_index, ArrowArray* output_array); + } // namespace iceberg diff --git a/src/iceberg/arrow_row_builder.cc b/src/iceberg/arrow_row_builder.cc index 91e3cd35e..7b3f2ac49 100644 --- a/src/iceberg/arrow_row_builder.cc +++ b/src/iceberg/arrow_row_builder.cc @@ -28,6 +28,8 @@ #include "iceberg/nanoarrow_status_internal.h" #include "iceberg/schema.h" #include "iceberg/schema_internal.h" +#include "iceberg/type.h" +#include "iceberg/util/decimal.h" namespace iceberg { @@ -137,6 +139,14 @@ Status AppendBytes(ArrowArray* array, std::span value) { return {}; } +Status AppendDecimal(ArrowArray* array, const Decimal& value, const DecimalType& type) { + ArrowDecimal decimal; + ArrowDecimalInit(&decimal, 128, type.precision(), type.scale()); + ArrowDecimalSetBytes(&decimal, value.native_endian_bytes()); + ICEBERG_NANOARROW_RETURN_UNEXPECTED(ArrowArrayAppendDecimal(array, &decimal)); + return {}; +} + Status AppendIntList(ArrowArray* array, const std::vector& values) { if (values.empty()) { return AppendNull(array); diff --git a/src/iceberg/arrow_row_builder_internal.h b/src/iceberg/arrow_row_builder_internal.h index e9c55f07a..436f017c5 100644 --- a/src/iceberg/arrow_row_builder_internal.h +++ b/src/iceberg/arrow_row_builder_internal.h @@ -131,6 +131,10 @@ ICEBERG_EXPORT Status AppendString(ArrowArray* array, std::string_view value); /// \brief Append a binary value to a nanoarrow array builder. ICEBERG_EXPORT Status AppendBytes(ArrowArray* array, std::span value); +/// \brief Append an unscaled decimal value to a Decimal128 array. +ICEBERG_EXPORT Status AppendDecimal(ArrowArray* array, const Decimal& value, + const DecimalType& type); + /// \brief Append a list of int32 values to a nanoarrow list array builder. ICEBERG_EXPORT Status AppendIntList(ArrowArray* array, const std::vector& values); diff --git a/src/iceberg/avro/avro_data_util.cc b/src/iceberg/avro/avro_data_util.cc index 297e46b52..6cf92f638 100644 --- a/src/iceberg/avro/avro_data_util.cc +++ b/src/iceberg/avro/avro_data_util.cc @@ -758,6 +758,7 @@ Status ExtractDatumFromArray(const ::arrow::Array& array, int64_t index, internal::checked_cast(array); std::string_view decimal_value = decimal_array.GetView(index); auto& fixed_datum = datum->value<::avro::GenericFixed>(); + decimal_value = decimal_value.substr(0, fixed_datum.schema()->fixedSize()); auto& bytes = fixed_datum.value(); bytes.assign(decimal_value.begin(), decimal_value.end()); std::ranges::reverse(bytes); diff --git a/src/iceberg/avro/avro_direct_encoder.cc b/src/iceberg/avro/avro_direct_encoder.cc index 5dcfd2511..027bfdb61 100644 --- a/src/iceberg/avro/avro_direct_encoder.cc +++ b/src/iceberg/avro/avro_direct_encoder.cc @@ -194,7 +194,9 @@ Status EncodeArrowToAvro(const ::avro::NodePtr& avro_node, ::avro::Encoder& enco if (avro_node->logicalType().type() == ::avro::LogicalType::DECIMAL) { const auto& decimal_array = internal::checked_cast(array); - std::string_view decimal_value = decimal_array.GetView(row_index); + // Avro fixed uses the schema's precision-dependent width, not all 16 Arrow bytes. + std::string_view decimal_value = + decimal_array.GetView(row_index).substr(0, avro_node->fixedSize()); ctx.bytes_scratch.assign(decimal_value.begin(), decimal_value.end()); // Arrow Decimal128 bytes are in little-endian order, Avro requires big-endian std::ranges::reverse(ctx.bytes_scratch); diff --git a/src/iceberg/deletes/position_delete_index.cc b/src/iceberg/deletes/position_delete_index.cc index 53e33e635..dd59cb4da 100644 --- a/src/iceberg/deletes/position_delete_index.cc +++ b/src/iceberg/deletes/position_delete_index.cc @@ -136,6 +136,10 @@ int64_t PositionDeleteIndex::Cardinality() const { return static_cast(bitmap_.Cardinality()); } +void PositionDeleteIndex::ForEach(const std::function& fn) const { + bitmap_.ForEach(fn); +} + void PositionDeleteIndex::Merge(const PositionDeleteIndex& other) { bitmap_.Or(other.bitmap_); delete_files_.insert(delete_files_.end(), other.delete_files_.begin(), diff --git a/src/iceberg/deletes/position_delete_index.h b/src/iceberg/deletes/position_delete_index.h index 6f301210e..828e43a2c 100644 --- a/src/iceberg/deletes/position_delete_index.h +++ b/src/iceberg/deletes/position_delete_index.h @@ -23,6 +23,7 @@ /// Index of deleted row positions for a data file. #include +#include #include #include #include @@ -66,6 +67,9 @@ class ICEBERG_EXPORT PositionDeleteIndex { /// \brief Get the number of deleted positions. int64_t Cardinality() const; + /// \brief Iterate over deleted positions in ascending order. + void ForEach(const std::function& fn) const; + /// \brief Merge another index into this one. /// \param other The index to merge (union operation) void Merge(const PositionDeleteIndex& other); diff --git a/src/iceberg/file_reader.h b/src/iceberg/file_reader.h index e08b2df49..0c514f859 100644 --- a/src/iceberg/file_reader.h +++ b/src/iceberg/file_reader.h @@ -57,6 +57,12 @@ class ICEBERG_EXPORT Reader { /// \brief Get the schema of the data. virtual Result Schema() = 0; + /// \brief Check whether a top-level field ID exists in the physical file schema. + /// \return Presence, or NotImplemented for readers without schema inspection. + virtual Result HasTopLevelField(int32_t field_id) { + return NotImplemented("Physical schema inspection is not supported"); + } + /// \brief Get the metadata of the file. virtual Result> Metadata() = 0; }; diff --git a/src/iceberg/inspect/metadata_table.h b/src/iceberg/inspect/metadata_table.h index 85000ee1e..bedd6cfa6 100644 --- a/src/iceberg/inspect/metadata_table.h +++ b/src/iceberg/inspect/metadata_table.h @@ -43,6 +43,7 @@ class ICEBERG_EXPORT MetadataTable { enum class Kind { kSnapshots, kHistory, + kPositionDeletes, }; /// \brief Maximum number of rows emitted in each Arrow batch. diff --git a/src/iceberg/inspect/position_deletes_table.cc b/src/iceberg/inspect/position_deletes_table.cc new file mode 100644 index 000000000..61d6855b4 --- /dev/null +++ b/src/iceberg/inspect/position_deletes_table.cc @@ -0,0 +1,767 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "iceberg/inspect/position_deletes_table.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include + +#include "iceberg/arrow_c_data_guard_internal.h" +#include "iceberg/arrow_c_data_util_internal.h" +#include "iceberg/arrow_row_builder_internal.h" +#include "iceberg/deletes/dv_util_internal.h" +#include "iceberg/deletes/position_delete_index.h" +#include "iceberg/expression/literal.h" +#include "iceberg/file_reader.h" +#include "iceberg/manifest/manifest_entry.h" +#include "iceberg/manifest/manifest_reader.h" +#include "iceberg/metadata_columns.h" +#include "iceberg/nanoarrow_status_internal.h" +#include "iceberg/partition_spec.h" +#include "iceberg/row/partition_values.h" +#include "iceberg/schema.h" +#include "iceberg/schema_field.h" +#include "iceberg/schema_internal.h" +#include "iceberg/snapshot.h" +#include "iceberg/table.h" +#include "iceberg/table_metadata.h" +#include "iceberg/transform.h" +#include "iceberg/type.h" +#include "iceberg/util/macros.h" +#include "iceberg/util/type_util.h" + +namespace iceberg { +namespace { + +using PartitionMapping = std::vector>; + +struct PositionDeletesSchema { + std::shared_ptr schema; + std::shared_ptr partition_type; + std::unordered_map partition_mappings; +}; + +struct UnifiedPartitionField { + int32_t partition_field_id; + int32_t source_id; + std::shared_ptr transform; + SchemaField field; +}; + +Status CollectFieldIds(const Type& type, std::unordered_set& field_ids) { + if (!type.is_nested()) { + return {}; + } + for (const auto& field : static_cast(type).fields()) { + if (!field_ids.insert(field.field_id()).second) { + return InvalidSchema("Duplicate field ID {} in position_deletes schema", + field.field_id()); + } + ICEBERG_RETURN_UNEXPECTED(CollectFieldIds(*field.type(), field_ids)); + } + return {}; +} + +Result TakeFreshFieldId(int64_t& candidate, + std::unordered_set& used_ids) { + constexpr int64_t kMaxFieldId = std::numeric_limits::max(); + for (int pass = 0; pass < 2; ++pass) { + while (candidate <= kMaxFieldId) { + const auto field_id = static_cast(candidate++); + if (used_ids.insert(field_id).second) { + return field_id; + } + } + candidate = 1; + } + return InvalidSchema("No field ID is available for position_deletes partition fields"); +} + +Result> ResolvePartitionField( + const PartitionField& partition_field, const Schema& current_schema) { + ICEBERG_ASSIGN_OR_RAISE(auto source_field, + current_schema.FindFieldById(partition_field.source_id())); + if (!source_field.has_value()) { + return std::nullopt; + } + return SchemaField::MakeOptional( + partition_field.field_id(), std::string(partition_field.name()), + partition_field.transform()->ResultType(source_field->get().type())); +} + +Result MakePositionDeletesSchema(const Table& table) { + ICEBERG_ASSIGN_OR_RAISE(auto table_schema, table.schema()); + + std::vector> specs = table.metadata()->partition_specs; + std::ranges::sort(specs, {}, [](const auto& spec) { + return spec == nullptr ? std::numeric_limits::min() : spec->spec_id(); + }); + std::ranges::reverse(specs); + + std::unordered_map fields_by_partition_id; + for (const auto& spec : specs) { + ICEBERG_PRECHECK(spec != nullptr, "Partition spec cannot be null"); + for (const auto& partition_field : spec->fields()) { + ICEBERG_PRECHECK( + partition_field.transform()->transform_type() != TransformType::kUnknown, + "Cannot build position_deletes partition type for unknown " + "transform on field {}", + partition_field.field_id()); + ICEBERG_ASSIGN_OR_RAISE(auto field, + ResolvePartitionField(partition_field, *table_schema)); + if (!field.has_value()) { + continue; + } + auto state = UnifiedPartitionField{ + .partition_field_id = partition_field.field_id(), + .source_id = partition_field.source_id(), + .transform = partition_field.transform(), + .field = std::move(*field), + }; + auto [existing_it, inserted] = + fields_by_partition_id.try_emplace(partition_field.field_id(), state); + if (inserted) { + continue; + } + + auto& existing = existing_it->second; + if (existing.source_id != partition_field.source_id()) { + return InvalidSchema("Partition field ID {} has conflicting source IDs {} and {}", + partition_field.field_id(), existing.source_id, + partition_field.source_id()); + } + const bool existing_void = + existing.transform->transform_type() == TransformType::kVoid; + const bool current_void = + partition_field.transform()->transform_type() == TransformType::kVoid; + if (!existing_void && !current_void && + *existing.transform != *partition_field.transform()) { + return InvalidSchema( + "Partition field ID {} has incompatible transforms {} and {}", + partition_field.field_id(), existing.transform->ToString(), + partition_field.transform()->ToString()); + } + if (!existing_void && !current_void && + *existing.field.type() != *state.field.type()) { + return InvalidSchema("Partition field ID {} has incompatible types {} and {}", + partition_field.field_id(), + existing.field.type()->ToString(), + state.field.type()->ToString()); + } + if (existing_void && !current_void) { + existing.transform = std::move(state.transform); + existing.field = existing.field.WithType(state.field.type()); + } + } + } + + std::vector unified_partition_fields; + unified_partition_fields.reserve(fields_by_partition_id.size()); + for (auto& [_, field] : fields_by_partition_id) { + unified_partition_fields.push_back(std::move(field)); + } + std::ranges::sort(unified_partition_fields, {}, + &UnifiedPartitionField::partition_field_id); + + constexpr std::array kMetadataFieldIds{ + MetadataColumns::kDeleteFilePathColumnId, + MetadataColumns::kDeleteFilePosColumnId, + MetadataColumns::kDeleteFileRowColumnId, + MetadataColumns::kPartitionColumnId, + MetadataColumns::kSpecIdColumnId, + MetadataColumns::kFilePathColumnId, + MetadataColumns::kContentOffsetColumnId, + MetadataColumns::kContentSizeInBytesColumnId, + }; + std::unordered_set used_ids(kMetadataFieldIds.begin(), + kMetadataFieldIds.end()); + std::unordered_set current_schema_ids; + ICEBERG_RETURN_UNEXPECTED(CollectFieldIds(*table_schema, current_schema_ids)); + for (int32_t field_id : current_schema_ids) { + if (!used_ids.insert(field_id).second) { + return InvalidSchema("Table field ID {} conflicts with a metadata field ID", + field_id); + } + } + for (const auto& historical_schema : table.metadata()->schemas) { + ICEBERG_PRECHECK(historical_schema != nullptr, "Table schema cannot be null"); + std::unordered_set historical_ids; + ICEBERG_RETURN_UNEXPECTED(CollectFieldIds(*historical_schema, historical_ids)); + used_ids.insert(historical_ids.begin(), historical_ids.end()); + } + + int64_t next_field_id = + std::max(1, static_cast(table.metadata()->last_column_id) + 1); + std::vector partition_fields; + partition_fields.reserve(unified_partition_fields.size()); + for (const auto& field : unified_partition_fields) { + ICEBERG_ASSIGN_OR_RAISE(auto field_id, TakeFreshFieldId(next_field_id, used_ids)); + partition_fields.push_back(SchemaField::MakeOptional( + field_id, field.field.name(), field.field.type(), field.field.doc())); + } + + auto partition_type = std::make_shared(partition_fields); + std::unordered_map partition_mappings; + for (const auto& spec : specs) { + PartitionMapping mapping(partition_fields.size()); + for (size_t output_idx = 0; output_idx < partition_fields.size(); ++output_idx) { + for (size_t input_idx = 0; input_idx < spec->fields().size(); ++input_idx) { + if (spec->fields()[input_idx].field_id() == + unified_partition_fields[output_idx].partition_field_id && + spec->fields()[input_idx].source_id() == + unified_partition_fields[output_idx].source_id) { + mapping[output_idx] = input_idx; + break; + } + } + } + partition_mappings.emplace(spec->spec_id(), std::move(mapping)); + } + + std::vector fields{ + MetadataColumns::kDeleteFilePath, + MetadataColumns::kDeleteFilePos, + SchemaField::MakeOptional( + MetadataColumns::kDeleteFileRowColumnId, + MetadataColumns::kDeleteFileRowFieldName, + std::make_shared(std::vector( + table_schema->fields().begin(), table_schema->fields().end())), + MetadataColumns::kDeleteFileRowDoc), + }; + if (!partition_fields.empty()) { + fields.push_back(SchemaField::MakeRequired( + MetadataColumns::kPartitionColumnId, "partition", partition_type, + "Partition that position delete row belongs to")); + } + fields.push_back( + SchemaField::MakeRequired(MetadataColumns::kSpecIdColumnId, "spec_id", int32(), + "Spec ID used to track the file containing a row")); + fields.push_back( + SchemaField::MakeRequired(MetadataColumns::kFilePathColumnId, "delete_file_path", + string(), "Path of the file in which a row is stored")); + + if (table.metadata()->format_version >= 3) { + fields.push_back(SchemaField::MakeOptional( + MetadataColumns::kContentOffsetColumnId, "content_offset", int64(), + "The offset in the DV where the content starts")); + fields.push_back(SchemaField::MakeOptional( + MetadataColumns::kContentSizeInBytesColumnId, "content_size_in_bytes", int64(), + "The length in bytes of the DV blob")); + } + + auto schema = std::make_shared(std::move(fields)); + ICEBERG_RETURN_UNEXPECTED(schema->HighestFieldId()); + return PositionDeletesSchema{ + .schema = std::move(schema), + .partition_type = std::move(partition_type), + .partition_mappings = std::move(partition_mappings), + }; +} + +template +Result GetLiteralValue(const Literal& literal, const Type& expected_type) { + if (const auto* value = std::get_if(&literal.value())) { + return value; + } + return InvalidArrowData("Partition value has type {} but metadata schema expects {}", + literal.type()->ToString(), expected_type.ToString()); +} + +Status AppendLiteralValue(ArrowArray* array, const Literal& literal, + const std::shared_ptr& type) { + if (literal.IsNull()) { + return AppendNull(array); + } + + const Literal* value = &literal; + std::optional promoted; + // Manifest readers may use a storage literal for a logical type (e.g. int for date). + // Apply schema promotions when needed; the typed extraction below checks storage. + if (*literal.type() != *type && IsPromotionAllowed(literal.type(), type)) { + ICEBERG_ASSIGN_OR_RAISE( + auto casted, literal.CastTo(std::static_pointer_cast(type))); + promoted.emplace(std::move(casted)); + value = &*promoted; + } + + switch (type->type_id()) { + case TypeId::kBoolean: { + ICEBERG_ASSIGN_OR_RAISE(auto bool_value, GetLiteralValue(*value, *type)); + return AppendBoolean(array, *bool_value); + } + case TypeId::kInt: + case TypeId::kDate: { + ICEBERG_ASSIGN_OR_RAISE(auto int_value, GetLiteralValue(*value, *type)); + return AppendInt(array, *int_value); + } + case TypeId::kLong: + case TypeId::kTime: + case TypeId::kTimestamp: + case TypeId::kTimestampTz: + case TypeId::kTimestampNs: + case TypeId::kTimestampTzNs: { + ICEBERG_ASSIGN_OR_RAISE(auto long_value, GetLiteralValue(*value, *type)); + return AppendInt(array, *long_value); + } + case TypeId::kFloat: { + ICEBERG_ASSIGN_OR_RAISE(auto float_value, GetLiteralValue(*value, *type)); + return AppendDouble(array, *float_value); + } + case TypeId::kDouble: { + ICEBERG_ASSIGN_OR_RAISE(auto double_value, GetLiteralValue(*value, *type)); + return AppendDouble(array, *double_value); + } + case TypeId::kString: { + ICEBERG_ASSIGN_OR_RAISE(auto string_value, + GetLiteralValue(*value, *type)); + return AppendString(array, *string_value); + } + case TypeId::kFixed: + case TypeId::kBinary: { + ICEBERG_ASSIGN_OR_RAISE(auto bytes_value, + GetLiteralValue>(*value, *type)); + return AppendBytes(array, *bytes_value); + } + case TypeId::kDecimal: { + ICEBERG_ASSIGN_OR_RAISE(auto decimal_value, + GetLiteralValue(*value, *type)); + return AppendDecimal(array, *decimal_value, static_cast(*type)); + } + case TypeId::kUuid: { + ICEBERG_ASSIGN_OR_RAISE(auto uuid_value, GetLiteralValue(*value, *type)); + return AppendBytes(array, uuid_value->bytes()); + } + default: + return InvalidArrowData("Unsupported partition type: {}", type->ToString()); + } +} + +Status AppendPartition(ArrowArray* array, const StructType& partition_type, + const PartitionValues& partition, + const PartitionMapping& mapping) { + ICEBERG_PRECHECK(std::cmp_equal(array->n_children, mapping.size()), + "Partition builder does not match metadata table schema"); + for (size_t output_idx = 0; output_idx < mapping.size(); ++output_idx) { + if (!mapping[output_idx].has_value()) { + ICEBERG_RETURN_UNEXPECTED(AppendNull(array->children[output_idx])); + continue; + } + const size_t input_idx = *mapping[output_idx]; + ICEBERG_PRECHECK(input_idx < partition.num_fields(), + "Partition values do not match partition spec"); + ICEBERG_ASSIGN_OR_RAISE(auto literal, partition.ValueAt(input_idx)); + ICEBERG_RETURN_UNEXPECTED( + AppendLiteralValue(array->children[output_idx], literal.get(), + partition_type.fields()[output_idx].type())); + } + ICEBERG_NANOARROW_RETURN_UNEXPECTED(ArrowArrayFinishElement(array)); + return {}; +} + +Status AppendOptionalInt(ArrowArray* array, const std::optional& value) { + return value.has_value() ? AppendInt(array, *value) : AppendNull(array); +} + +struct OutputLayout { + explicit OutputLayout(bool has_partition, bool is_v3) + : partition(has_partition ? std::optional(3) : std::nullopt), + spec_id(has_partition ? 4 : 3), + delete_file_path(spec_id + 1), + content_offset(is_v3 ? std::optional(delete_file_path + 1) + : std::nullopt), + content_size(is_v3 ? std::optional(delete_file_path + 2) : std::nullopt) { + } + + std::optional partition; + size_t spec_id; + size_t delete_file_path; + std::optional content_offset; + std::optional content_size; +}; + +struct DeletedRowValue { + const ArrowSchema* schema; + const ArrowArray* array; + const ArrowArrayView* view; + int64_t row_index; +}; + +// Bound emitted batches while retaining the existing materialized inspection model. +class PositionDeleteBatchBuilder { + public: + PositionDeleteBatchBuilder(const Schema& schema, ArrowRowBuilder builder) + : schema_(schema), builder_(std::move(builder)) {} + ~PositionDeleteBatchBuilder() { + for (auto& batch : batches_) { + if (batch.release) ArrowArrayRelease(&batch); + } + } + ArrowArray* column(int64_t index) { return builder_.column(index); } + Status FinishRow() { + ICEBERG_RETURN_UNEXPECTED(builder_.FinishRow()); + if (builder_.num_rows() == MetadataTable::kBatchSize) { + ICEBERG_ASSIGN_OR_RAISE(auto batch, std::move(builder_).Finish()); + batches_.push_back(batch); + ICEBERG_ASSIGN_OR_RAISE(builder_, ArrowRowBuilder::Make(schema_)); + } + return {}; + } + Result Finish(const Schema& projected_schema, + ProjectionContext* projection) { + if (builder_.num_rows() != 0 || batches_.empty()) { + ICEBERG_ASSIGN_OR_RAISE(auto batch, std::move(builder_).Finish()); + batches_.push_back(batch); + } + ArrowSchema schema{}; + ICEBERG_RETURN_UNEXPECTED(ToArrowSchema(projected_schema, &schema)); + internal::ArrowSchemaGuard schema_guard(&schema); + ArrowArrayStream stream{}; + ICEBERG_NANOARROW_RETURN_UNEXPECTED( + ArrowBasicArrayStreamInit(&stream, &schema, batches_.size())); + nanoarrow::UniqueArrayStream stream_guard; + ArrowArrayStreamMove(&stream, stream_guard.get()); + for (size_t i = 0; i < batches_.size(); ++i) { + ArrowArray batch{}; + ArrowArrayMove(&batches_[i], &batch); + if (projection != nullptr) { + std::vector rows(static_cast(batch.length)); + std::iota(rows.begin(), rows.end(), 0); + ICEBERG_ASSIGN_OR_RAISE(batch, ProjectBatch(&batch, rows, *projection)); + } + ArrowBasicArrayStreamSetArray(stream_guard.get(), i, &batch); + } + ArrowArrayStreamMove(stream_guard.get(), &stream); + return stream; + } + + private: + const Schema& schema_; + ArrowRowBuilder builder_; + std::vector batches_; +}; + +Status AppendPositionDeleteRow( + PositionDeleteBatchBuilder& builder, const OutputLayout& layout, + const DataFile& delete_file, std::string_view data_file_path, int64_t pos, + int32_t spec_id, const StructType& partition_type, + const std::unordered_map& partition_mappings, + const DeletedRowValue* deleted_row) { + ICEBERG_RETURN_UNEXPECTED(AppendString(builder.column(0), data_file_path)); + ICEBERG_RETURN_UNEXPECTED(AppendInt(builder.column(1), pos)); + if (deleted_row == nullptr) { + ICEBERG_RETURN_UNEXPECTED(AppendNull(builder.column(2))); + } else { + ICEBERG_RETURN_UNEXPECTED(AppendArrayValue(*deleted_row->schema, *deleted_row->array, + *deleted_row->view, deleted_row->row_index, + builder.column(2))); + } + + if (layout.partition.has_value()) { + auto mapping = partition_mappings.find(spec_id); + ICEBERG_PRECHECK(mapping != partition_mappings.end(), + "Partition spec ID {} not found", spec_id); + ICEBERG_RETURN_UNEXPECTED(AppendPartition(builder.column(*layout.partition), + partition_type, delete_file.partition, + mapping->second)); + } + + ICEBERG_RETURN_UNEXPECTED(AppendInt(builder.column(layout.spec_id), spec_id)); + ICEBERG_RETURN_UNEXPECTED( + AppendString(builder.column(layout.delete_file_path), delete_file.file_path)); + if (layout.content_offset.has_value()) { + ICEBERG_RETURN_UNEXPECTED(AppendOptionalInt(builder.column(*layout.content_offset), + delete_file.content_offset)); + ICEBERG_RETURN_UNEXPECTED(AppendOptionalInt(builder.column(*layout.content_size), + delete_file.content_size_in_bytes)); + } + return builder.FinishRow(); +} + +SchemaField ClearPositionDeleteReadDefaults(SchemaField field) { + return field.WithInitialDefault(nullptr).WithWriteDefault(nullptr); +} + +std::shared_ptr MakePositionDeleteReadType(const std::shared_ptr& type) { + switch (type->type_id()) { + case TypeId::kStruct: { + const auto& struct_type = static_cast(*type); + std::vector fields; + fields.reserve(struct_type.fields().size()); + for (const auto& field : struct_type.fields()) { + fields.push_back(ClearPositionDeleteReadDefaults( + field.WithType(MakePositionDeleteReadType(field.type()))) + .AsOptional()); + } + return std::make_shared(std::move(fields)); + } + case TypeId::kList: { + const auto& list_type = static_cast(*type); + return std::make_shared( + ClearPositionDeleteReadDefaults(list_type.element().WithType( + MakePositionDeleteReadType(list_type.element().type())))); + } + case TypeId::kMap: { + const auto& map_type = static_cast(*type); + return std::make_shared( + ClearPositionDeleteReadDefaults(map_type.key()), + ClearPositionDeleteReadDefaults(map_type.value().WithType( + MakePositionDeleteReadType(map_type.value().type())))); + } + default: + return type; + } +} + +std::shared_ptr PositionDeleteFileSchema(const Schema& table_schema) { + std::vector fields{ + MetadataColumns::kDeleteFilePath, + MetadataColumns::kDeleteFilePos, + }; + std::vector row_fields; + row_fields.reserve(table_schema.fields().size()); + for (const auto& field : table_schema.fields()) { + row_fields.push_back(ClearPositionDeleteReadDefaults( + field.WithType(MakePositionDeleteReadType(field.type()))) + .AsOptional()); + } + fields.push_back(SchemaField::MakeOptional( + MetadataColumns::kDeleteFileRowColumnId, MetadataColumns::kDeleteFileRowFieldName, + std::make_shared(std::move(row_fields)), + MetadataColumns::kDeleteFileRowDoc)); + return std::make_shared(std::move(fields)); +} + +struct PositionDeleteReader { + std::unique_ptr reader; + bool has_row; +}; + +Result OpenPositionDeleteFile(const DataFile& file, + const Schema& table_schema, + const std::shared_ptr& io) { + ICEBERG_PRECHECK(file.file_format == FileFormatType::kParquet, + "Unsupported position delete format: {}", ToString(file.file_format)); + ReaderOptions options{.path = file.file_path, + .length = static_cast(file.file_size_in_bytes), + .io = io, + .projection = PositionDeleteFileSchema(table_schema)}; + ICEBERG_ASSIGN_OR_RAISE(auto reader, + ReaderFactoryRegistry::Open(file.file_format, options)); + ICEBERG_ASSIGN_OR_RAISE( + auto has_row, reader->HasTopLevelField(MetadataColumns::kDeleteFileRowColumnId)); + return PositionDeleteReader{.reader = std::move(reader), .has_row = has_row}; +} + +Status AppendParquetDeletes( + PositionDeleteBatchBuilder& builder, const OutputLayout& layout, + const std::shared_ptr& delete_file, int32_t spec_id, + const StructType& partition_type, + const std::unordered_map& partition_mappings, + const Schema& table_schema, const std::shared_ptr& io) { + ICEBERG_ASSIGN_OR_RAISE(auto delete_reader, + OpenPositionDeleteFile(*delete_file, table_schema, io)); + auto& reader = delete_reader.reader; + ICEBERG_ASSIGN_OR_RAISE(auto arrow_schema, reader->Schema()); + internal::ArrowSchemaGuard schema_guard(&arrow_schema); + constexpr int64_t expected_columns = 3; + ICEBERG_PRECHECK(arrow_schema.n_children == expected_columns, + "Position delete reader returned {} columns, expected {}", + arrow_schema.n_children, expected_columns); + + ArrowArrayView array_view; + internal::ArrowArrayViewGuard view_guard(&array_view); + ArrowError error; + ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( + ArrowArrayViewInitFromSchema(&array_view, &arrow_schema, &error), error); + + while (true) { + ICEBERG_ASSIGN_OR_RAISE(auto batch_opt, reader->Next()); + if (!batch_opt.has_value()) { + break; + } + + auto& batch = *batch_opt; + internal::ArrowArrayGuard batch_guard(&batch); + ICEBERG_NANOARROW_RETURN_UNEXPECTED_WITH_ERROR( + ArrowArrayViewSetArray(&array_view, &batch, &error), error); + ICEBERG_PRECHECK(batch.n_children == expected_columns, + "Position delete batch has {} columns, expected {}", + batch.n_children, expected_columns); + + const auto* path_view = array_view.children[0]; + const auto* pos_view = array_view.children[1]; + if (ArrowArrayViewComputeNullCount(path_view) != 0 || + ArrowArrayViewComputeNullCount(pos_view) != 0) { + return InvalidArrowData( + "position delete file has null values in required pos/file_path columns"); + } + if (delete_reader.has_row && + ArrowArrayViewComputeNullCount(array_view.children[2]) != 0) { + return InvalidArrowData("position delete file has null row values"); + } + + const int64_t* positions = pos_view->buffer_views[1].data.as_int64 + pos_view->offset; + for (int64_t row = 0; row < batch.length; ++row) { + ICEBERG_PRECHECK(positions[row] >= 0, "Invalid negative delete position: {}", + positions[row]); + const ArrowStringView path = ArrowArrayViewGetStringUnsafe(path_view, row); + std::optional deleted_row; + if (delete_reader.has_row) { + deleted_row.emplace(DeletedRowValue{ + .schema = arrow_schema.children[2], + .array = batch.children[2], + .view = array_view.children[2], + .row_index = row, + }); + } + ICEBERG_RETURN_UNEXPECTED(AppendPositionDeleteRow( + builder, layout, *delete_file, + std::string_view(path.data, static_cast(path.size_bytes)), + positions[row], spec_id, partition_type, partition_mappings, + deleted_row ? &*deleted_row : nullptr)); + } + } + + return reader->Close(); +} + +Status AppendDeletionVector( + PositionDeleteBatchBuilder& builder, const OutputLayout& layout, + const std::shared_ptr& delete_file, int32_t spec_id, + const StructType& partition_type, + const std::unordered_map& partition_mappings, + const std::shared_ptr& io) { + ICEBERG_PRECHECK(delete_file->referenced_data_file.has_value(), + "Deletion vector requires a referenced data file"); + ICEBERG_ASSIGN_OR_RAISE(auto positions, DVUtil::ReadDV(delete_file, io)); + + Status status = {}; + positions.ForEach([&](int64_t pos) { + if (status.has_value()) { + status = AppendPositionDeleteRow(builder, layout, *delete_file, + *delete_file->referenced_data_file, pos, spec_id, + partition_type, partition_mappings, nullptr); + } + }); + return status; +} + +} // namespace + +PositionDeletesTable::PositionDeletesTable( + std::shared_ptr table, std::shared_ptr schema, + std::shared_ptr partition_type, + std::unordered_map partition_mappings) + : MetadataTable(std::move(table)), + schema_(std::move(schema)), + source_metadata_(source_table()->metadata()), + partition_type_(std::move(partition_type)), + partition_mappings_(std::move(partition_mappings)) {} + +PositionDeletesTable::~PositionDeletesTable() = default; + +Result> PositionDeletesTable::Make( + std::shared_ptr
table) { + ICEBERG_PRECHECK(table != nullptr, "Table cannot be null"); + ICEBERG_ASSIGN_OR_RAISE(auto schema, MakePositionDeletesSchema(*table)); + return std::unique_ptr(new PositionDeletesTable( + std::move(table), std::move(schema.schema), std::move(schema.partition_type), + std::move(schema.partition_mappings))); +} + +Result PositionDeletesTable::Scan() { return Scan(*schema()); } + +Result PositionDeletesTable::Scan(const Schema& projected_schema) { + const auto& current = *source_table()->metadata(); + ICEBERG_PRECHECK( + current.format_version == source_metadata_->format_version && + current.current_schema_id == source_metadata_->current_schema_id && + current.default_spec_id == source_metadata_->default_spec_id && + std::ranges::equal( + current.partition_specs, source_metadata_->partition_specs, + [](const auto& lhs, const auto& rhs) { return *lhs == *rhs; }), + "Source schema or partitioning changed; recreate the position_deletes table"); + std::optional projection; + if (*schema() != projected_schema) { + ICEBERG_ASSIGN_OR_RAISE( + auto context, + ProjectionContext::Make(*schema(), projected_schema, + ProjectionContext::ResolveProjectBatchFunction())); + projection.emplace(std::move(context)); + } + ICEBERG_ASSIGN_OR_RAISE(auto row_builder, ArrowRowBuilder::Make(*schema())); + PositionDeleteBatchBuilder builder(*schema(), std::move(row_builder)); + const bool has_partition = !partition_type_->fields().empty(); + const bool is_v3 = source_table()->metadata()->format_version >= 3; + const OutputLayout layout(has_partition, is_v3); + + if (source_table()->metadata()->current_snapshot_id != kInvalidSnapshotId) { + ICEBERG_ASSIGN_OR_RAISE(auto current_snapshot, source_table()->current_snapshot()); + SnapshotReader cache(current_snapshot.get()); + ICEBERG_ASSIGN_OR_RAISE(auto manifests, cache.DeleteManifests(source_table()->io())); + ICEBERG_ASSIGN_OR_RAISE(auto table_schema, source_table()->schema()); + ICEBERG_ASSIGN_OR_RAISE(auto specs_ref, source_table()->specs()); + const auto& specs = specs_ref.get(); + + for (const auto& manifest : manifests) { + ICEBERG_ASSIGN_OR_RAISE( + auto reader, + ManifestReader::Make(manifest, source_table()->io(), table_schema, specs)); + ICEBERG_ASSIGN_OR_RAISE(auto entries, reader->LiveEntries()); + for (const auto& entry : entries) { + ICEBERG_PRECHECK(entry.data_file != nullptr, + "Manifest entry must have a data file"); + const auto& file = entry.data_file; + if (file->content != DataFile::Content::kPositionDeletes) { + continue; + } + const int32_t spec_id = + file->partition_spec_id.value_or(manifest.partition_spec_id); + if (file->IsDeletionVector()) { + ICEBERG_RETURN_UNEXPECTED( + AppendDeletionVector(builder, layout, file, spec_id, *partition_type_, + partition_mappings_, source_table()->io())); + } else { + ICEBERG_RETURN_UNEXPECTED(AppendParquetDeletes( + builder, layout, file, spec_id, *partition_type_, partition_mappings_, + *table_schema, source_table()->io())); + } + } + } + } + + return builder.Finish(projected_schema, projection ? &*projection : nullptr); +} + +} // namespace iceberg diff --git a/src/iceberg/inspect/position_deletes_table.h b/src/iceberg/inspect/position_deletes_table.h new file mode 100644 index 000000000..0f54bf216 --- /dev/null +++ b/src/iceberg/inspect/position_deletes_table.h @@ -0,0 +1,72 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#pragma once + +/// \file iceberg/inspect/position_deletes_table.h +/// \brief Define the position_deletes metadata table. + +#include +#include +#include +#include +#include + +#include "iceberg/iceberg_export.h" +#include "iceberg/inspect/metadata_table.h" +#include "iceberg/result.h" +#include "iceberg/type_fwd.h" + +namespace iceberg { + +/// \brief Metadata table that expands physical position-delete storage into rows. +class ICEBERG_EXPORT PositionDeletesTable : public MetadataTable { + public: + /// \brief Create a position_deletes metadata table for a source table. + /// \param table Source table whose current snapshot will be inspected. + /// \return A metadata table or an error if its schema cannot be constructed. + static Result> Make(std::shared_ptr
table); + + ~PositionDeletesTable() override; + + Kind kind() const noexcept override { return Kind::kPositionDeletes; } + + const std::shared_ptr& schema() const override { return schema_; } + + /// \brief Materialize current position deletes as an Arrow stream. + /// Recreate this metadata table if the source schema, format, or partition specs + /// change. + Result Scan() override; + + /// \brief Scan with a projection of complete top-level fields. + Result Scan(const Schema& projected_schema); + + private: + PositionDeletesTable( + std::shared_ptr
table, std::shared_ptr schema, + std::shared_ptr partition_type, + std::unordered_map>> partition_mappings); + + std::shared_ptr schema_; + std::shared_ptr source_metadata_; + std::shared_ptr partition_type_; + std::unordered_map>> partition_mappings_; +}; + +} // namespace iceberg diff --git a/src/iceberg/manifest/manifest_adapter.cc b/src/iceberg/manifest/manifest_adapter.cc index b26249a6c..8b95cf525 100644 --- a/src/iceberg/manifest/manifest_adapter.cc +++ b/src/iceberg/manifest/manifest_adapter.cc @@ -147,8 +147,9 @@ Status ManifestEntryAdapter::AppendPartitionValues( AppendInt(child_array, std::get(partition_value.value()))); break; case TypeId::kDecimal: - ICEBERG_RETURN_UNEXPECTED(AppendBytes( - child_array, std::get(partition_value.value()).ToBytes())); + ICEBERG_RETURN_UNEXPECTED( + AppendDecimal(child_array, std::get(partition_value.value()), + static_cast(*partition_field.type()))); break; case TypeId::kUuid: ICEBERG_RETURN_UNEXPECTED( diff --git a/src/iceberg/manifest/manifest_reader.cc b/src/iceberg/manifest/manifest_reader.cc index 7c920b048..773f05053 100644 --- a/src/iceberg/manifest/manifest_reader.cc +++ b/src/iceberg/manifest/manifest_reader.cc @@ -45,6 +45,7 @@ #include "iceberg/type.h" #include "iceberg/util/checked_cast.h" #include "iceberg/util/content_file_util.h" +#include "iceberg/util/decimal.h" #include "iceberg/util/macros.h" #include "nanoarrow/common/inline_types.h" @@ -416,6 +417,17 @@ Status ParsePartitionValues(ArrowArrayView* view, int64_t row_idx, case ArrowType::NANOARROW_TYPE_DOUBLE: partition.AddValue(Literal::Double(ArrowArrayViewGetDoubleUnsafe(view, row_idx))); break; + case ArrowType::NANOARROW_TYPE_DECIMAL128: { + const auto& type = internal::checked_cast(*field_type); + ArrowDecimal value; + ArrowDecimalInit(&value, 128, type.precision(), type.scale()); + ArrowArrayViewGetDecimalUnsafe(view, row_idx, &value); + Decimal unscaled(static_cast(value.words[value.high_word_index]), + value.words[value.low_word_index]); + partition.AddValue( + Literal::Decimal(unscaled.value(), type.precision(), type.scale())); + break; + } case ArrowType::NANOARROW_TYPE_STRING: { auto str_value = ArrowArrayViewGetStringUnsafe(view, row_idx); partition.AddValue( diff --git a/src/iceberg/parquet/parquet_reader.cc b/src/iceberg/parquet/parquet_reader.cc index 1f8146107..b98825c4e 100644 --- a/src/iceberg/parquet/parquet_reader.cc +++ b/src/iceberg/parquet/parquet_reader.cc @@ -365,6 +365,15 @@ class ParquetReader::Impl { return arrow_schema; } + Result HasTopLevelField(int32_t field_id) { + if (!reader_) return Invalid("Reader is not opened"); + const auto* root = reader_->parquet_reader()->metadata()->schema()->group_node(); + for (int i = 0; i < root->field_count(); ++i) { + if (root->field(i)->field_id() == field_id) return true; + } + return false; + } + Result> Metadata() { if (reader_ == nullptr) { return Invalid("Reader is not opened"); @@ -472,6 +481,10 @@ Result> ParquetReader::Next() { return impl_->Next(); Result ParquetReader::Schema() { return impl_->Schema(); } +Result ParquetReader::HasTopLevelField(int32_t field_id) { + return impl_->HasTopLevelField(field_id); +} + Result> ParquetReader::Metadata() { return impl_->Metadata(); } diff --git a/src/iceberg/parquet/parquet_reader.h b/src/iceberg/parquet/parquet_reader.h index 3a8e57cf3..c8ec52a06 100644 --- a/src/iceberg/parquet/parquet_reader.h +++ b/src/iceberg/parquet/parquet_reader.h @@ -42,6 +42,8 @@ class ICEBERG_BUNDLE_EXPORT ParquetReader : public Reader { Result Schema() final; + Result HasTopLevelField(int32_t field_id) final; + Result> Metadata() final; private: diff --git a/src/iceberg/test/CMakeLists.txt b/src/iceberg/test/CMakeLists.txt index 6f6ff7603..2e15051b6 100644 --- a/src/iceberg/test/CMakeLists.txt +++ b/src/iceberg/test/CMakeLists.txt @@ -194,6 +194,9 @@ if(ICEBERG_BUILD_BUNDLE) metadata_table_test.cc snapshots_table_test.cc) + add_iceberg_test(position_deletes_table_test USE_BUNDLE SOURCES + position_deletes_table_test.cc) + add_iceberg_test(eval_expr_test USE_BUNDLE SOURCES diff --git a/src/iceberg/test/avro_data_test.cc b/src/iceberg/test/avro_data_test.cc index 4ea7855ce..d86bfe436 100644 --- a/src/iceberg/test/avro_data_test.cc +++ b/src/iceberg/test/avro_data_test.cc @@ -1132,20 +1132,19 @@ const std::vector kExtractDatumTestCases = { { .name = "Decimal", .iceberg_type = decimal(10, 2), - .arrow_json = R"([{"a": "0.00"}, {"a": "10.01"}, {"a": "20.02"}])", + .arrow_json = R"([{"a": "-10.01"}, {"a": "0.00"}, {"a": "10.01"}])", .value_verifier = [](const ::avro::GenericDatum& datum, int i) { const auto& record = datum.value<::avro::GenericRecord>(); const auto& fixed = record.fieldAt(0).value<::avro::GenericFixed>(); const auto& bytes = fixed.value(); + ASSERT_EQ(bytes.size(), 5); auto decimal = ::arrow::Decimal128::FromBigEndian( reinterpret_cast(bytes.data()), bytes.size()) .ValueOrDie(); - int64_t expected_unscaled = i * 1000 + i; - EXPECT_EQ(decimal.low_bits(), static_cast(expected_unscaled)); - EXPECT_EQ(decimal.high_bits(), 0); + EXPECT_EQ(decimal, ::arrow::Decimal128((i - 1) * 1001)); }, }, { diff --git a/src/iceberg/test/position_deletes_table_test.cc b/src/iceberg/test/position_deletes_table_test.cc new file mode 100644 index 000000000..860f724e0 --- /dev/null +++ b/src/iceberg/test/position_deletes_table_test.cc @@ -0,0 +1,701 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "iceberg/inspect/position_deletes_table.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include + +#include "iceberg/arrow/arrow_register.h" +#include "iceberg/arrow/arrow_status_internal.h" +#include "iceberg/data/position_delete_writer.h" +#include "iceberg/deletes/dv_writer.h" +#include "iceberg/file_format.h" +#include "iceberg/file_writer.h" +#include "iceberg/inspect/metadata_table.h" +#include "iceberg/manifest/manifest_entry.h" +#include "iceberg/metadata_columns.h" +#include "iceberg/parquet/parquet_register.h" +#include "iceberg/partition_spec.h" +#include "iceberg/row/partition_values.h" +#include "iceberg/schema.h" +#include "iceberg/schema_internal.h" +#include "iceberg/snapshot.h" +#include "iceberg/table.h" +#include "iceberg/test/matchers.h" +#include "iceberg/test/mock_catalog.h" +#include "iceberg/test/scan_test_base.h" +#include "iceberg/util/macros.h" + +namespace iceberg { +namespace { + +using ::testing::ElementsAre; + +class PositionDeletesTableTest : public ScanTestBase { + protected: + void SetUp() override { + ScanTestBase::SetUp(); + parquet::RegisterAll(); + arrow::RegisterAll(); + } + + std::shared_ptr WritePositionDeletes( + std::string path, const std::vector>& deletes, + const std::shared_ptr& spec, const PartitionValues& partition) { + PositionDeleteWriterOptions options{ + .path = std::move(path), + .schema = schema_, + .spec = spec, + .partition = partition, + .format = FileFormatType::kParquet, + .io = file_io_, + .flush_threshold = 10000, + .properties = {{"write.parquet.compression-codec", "uncompressed"}}, + }; + auto writer = PositionDeleteWriter::Make(options).value(); + for (const auto& [file_path, pos] : deletes) { + ICEBERG_THROW_NOT_OK(writer->WriteDelete(file_path, pos)); + } + ICEBERG_THROW_NOT_OK(writer->Close()); + return writer->Metadata().value().data_files.front(); + } + + Result> WritePositionDeletesWithRows( + std::string path, std::string_view json, std::vector row_fields, + bool row_required, const std::shared_ptr& spec, + const PartitionValues& partition) { + auto row_type = std::make_shared(std::move(row_fields)); + auto row_field = + row_required ? SchemaField::MakeRequired(MetadataColumns::kDeleteFileRowColumnId, + MetadataColumns::kDeleteFileRowFieldName, + std::move(row_type), + MetadataColumns::kDeleteFileRowDoc) + : SchemaField::MakeOptional(MetadataColumns::kDeleteFileRowColumnId, + MetadataColumns::kDeleteFileRowFieldName, + std::move(row_type), + MetadataColumns::kDeleteFileRowDoc); + auto delete_schema = std::make_shared(std::vector{ + MetadataColumns::kDeleteFilePath, + MetadataColumns::kDeleteFilePos, + std::move(row_field), + }); + ICEBERG_ASSIGN_OR_RAISE( + auto writer, WriterFactoryRegistry::Open( + FileFormatType::kParquet, + WriterOptions{ + .path = path, + .schema = delete_schema, + .io = file_io_, + .properties = WriterProperties::FromMap( + {{"write.parquet.compression-codec", "uncompressed"}}), + })); + + ArrowSchema arrow_c_schema; + ICEBERG_THROW_NOT_OK(ToArrowSchema(*delete_schema, &arrow_c_schema)); + auto arrow_type = ::arrow::ImportType(&arrow_c_schema).ValueOrDie(); + auto rows = ::arrow::json::ArrayFromJSONString(::arrow::struct_(arrow_type->fields()), + std::string(json)) + .ValueOrDie(); + ArrowArray arrow_array; + ICEBERG_ARROW_RETURN_NOT_OK(::arrow::ExportArray(*rows, &arrow_array)); + ICEBERG_RETURN_UNEXPECTED(writer->Write(&arrow_array)); + ICEBERG_RETURN_UNEXPECTED(writer->Close()); + ICEBERG_ASSIGN_OR_RAISE(auto length, writer->length()); + + return std::make_shared(DataFile{ + .content = DataFile::Content::kPositionDeletes, + .file_path = std::move(path), + .file_format = FileFormatType::kParquet, + .partition = partition, + .record_count = rows->length(), + .file_size_in_bytes = length, + .partition_spec_id = spec->spec_id(), + }); + } + + std::vector> WriteDeletionVectors( + std::string path, + const std::vector>>& deletes, + const std::shared_ptr& spec, const PartitionValues& partition) { + auto writer = DVWriter::Make(DVWriterOptions{ + .path = std::move(path), + .io = file_io_, + .load_previous_deletes = [](std::string_view) + -> Result> { + return std::nullopt; + }, + }) + .value(); + for (const auto& [file_path, positions] : deletes) { + for (int64_t pos : positions) { + ICEBERG_THROW_NOT_OK(writer->Delete(file_path, pos, spec, partition)); + } + } + ICEBERG_THROW_NOT_OK(writer->Close()); + return writer->Metadata().value().data_files; + } + + std::shared_ptr MakeSnapshot(int8_t table_format_version, int64_t snapshot_id, + int64_t sequence_number, + const std::vector& manifests) { + auto manifest_list = WriteManifestList(table_format_version, snapshot_id, 0, + sequence_number, manifests); + return std::make_shared(Snapshot{ + .snapshot_id = snapshot_id, + .parent_snapshot_id = std::nullopt, + .sequence_number = sequence_number, + .timestamp_ms = TimePointMsFromUnixMs(1609459200000L), + .manifest_list = std::move(manifest_list), + .summary = {{"operation", "delete"}}, + .schema_id = schema_->schema_id(), + }); + } + + std::shared_ptr
MakeTable(int8_t format_version, + const std::shared_ptr& spec, + std::shared_ptr snapshot = nullptr) { + return MakeTableWithSpecs(format_version, {spec}, spec, std::move(snapshot)); + } + + std::shared_ptr
MakeTableWithSpecs( + int8_t format_version, std::vector> specs, + const std::shared_ptr& default_spec, + std::shared_ptr snapshot = nullptr) { + std::vector> snapshots; + if (snapshot != nullptr) { + snapshots.push_back(snapshot); + } + int32_t last_partition_id = PartitionSpec::kInvalidPartitionFieldId; + for (const auto& spec : specs) { + last_partition_id = std::max(last_partition_id, spec->last_assigned_field_id()); + } + auto metadata = std::make_shared(TableMetadata{ + .format_version = format_version, + .table_uuid = "test-table-uuid", + .location = "/tmp/table", + .last_sequence_number = snapshot == nullptr ? 0 : snapshot->sequence_number, + .last_updated_ms = TimePointMsFromUnixMs(1609459200000L), + .last_column_id = schema_->HighestFieldId().value(), + .schemas = {schema_}, + .current_schema_id = schema_->schema_id(), + .partition_specs = std::move(specs), + .default_spec_id = default_spec->spec_id(), + .last_partition_id = last_partition_id, + .current_snapshot_id = + snapshot == nullptr ? kInvalidSnapshotId : snapshot->snapshot_id, + .snapshots = std::move(snapshots), + .default_sort_order_id = 0, + }); + return Table::Make(TableIdentifier{.name = "table"}, std::move(metadata), + "/tmp/table/metadata.json", file_io_, + std::make_shared<::testing::NiceMock>()) + .value(); + } + + std::unique_ptr MakePositionDeletesTable( + const std::shared_ptr
& table) { + return MetadataTable::Make(table).value(); + } + + static std::shared_ptr<::arrow::RecordBatch> Import(ArrowArrayStream stream, + const Schema& schema) { + auto reader = ::arrow::ImportRecordBatchReader(&stream).ValueOrDie(); + auto batch = reader->Next().ValueOrDie(); + EXPECT_NE(batch, nullptr); + EXPECT_EQ(batch->num_columns(), schema.fields().size()); + return batch; + } +}; + +TEST_P(PositionDeletesTableTest, SchemaAndEmptyScan) { + auto table = MakePositionDeletesTable(MakeTable(GetParam(), unpartitioned_spec_)); + + EXPECT_EQ(table->kind(), MetadataTable::Kind::kPositionDeletes); + EXPECT_EQ(table->source_table()->name().name, "table"); + + std::vector names; + for (const auto& field : table->schema()->fields()) { + names.emplace_back(field.name()); + } + std::vector expected{"file_path", "pos", "row", "spec_id", + "delete_file_path"}; + if (GetParam() >= 3) { + expected.push_back("content_offset"); + expected.push_back("content_size_in_bytes"); + } + EXPECT_EQ(names, expected); + + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan()); + auto batch = Import(std::move(array), *table->schema()); + EXPECT_EQ(batch->num_rows(), 0); +} + +TEST_P(PositionDeletesTableTest, ExpandsParquetAndProjectsColumns) { + auto delete_file = WritePositionDeletes( + "position-deletes.parquet", {{"data-a.parquet", 2}, {"data-b.parquet", 7}}, + partitioned_spec_, PartitionValues(Literal::Int(11))); + constexpr int64_t kSnapshotId = 10; + auto manifest = WriteDeleteManifest( + GetParam(), kSnapshotId, + {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, delete_file)}, + partitioned_spec_); + auto table = MakePositionDeletesTable( + MakeTable(GetParam(), partitioned_spec_, + MakeSnapshot(GetParam(), kSnapshotId, 1, {manifest}))); + + const auto& fields = table->schema()->fields(); + auto projected = + std::make_unique(std::vector{fields[5], fields[1], fields[3]}); + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan(*projected)); + auto batch = Import(std::move(array), *projected); + + ASSERT_EQ(batch->num_rows(), 2); + EXPECT_EQ(batch->schema()->field(0)->name(), "delete_file_path"); + EXPECT_EQ(batch->schema()->field(1)->name(), "pos"); + EXPECT_EQ(batch->schema()->field(2)->name(), "partition"); + auto delete_paths = std::static_pointer_cast<::arrow::StringArray>(batch->column(0)); + auto positions = std::static_pointer_cast<::arrow::Int64Array>(batch->column(1)); + auto partitions = std::static_pointer_cast<::arrow::StructArray>(batch->column(2)); + auto buckets = std::static_pointer_cast<::arrow::Int32Array>(partitions->field(0)); + std::vector actual_delete_paths{delete_paths->GetString(0), + delete_paths->GetString(1)}; + std::vector actual_positions{positions->Value(0), positions->Value(1)}; + std::vector actual_buckets{buckets->Value(0), buckets->Value(1)}; + EXPECT_THAT(actual_delete_paths, + ElementsAre("position-deletes.parquet", "position-deletes.parquet")); + EXPECT_THAT(actual_positions, ElementsAre(2, 7)); + EXPECT_THAT(actual_buckets, ElementsAre(11, 11)); +} + +TEST_F(PositionDeletesTableTest, SurfacesDeletedRowsFromParquet) { + ICEBERG_UNWRAP_OR_FAIL( + auto delete_file, + WritePositionDeletesWithRows( + "position-deletes-with-rows.parquet", + R"([["data.parquet", 5, [42, "deleted"]]])", + std::vector(schema_->fields().begin(), schema_->fields().end()), + true, unpartitioned_spec_, PartitionValues{})); + constexpr int64_t kSnapshotId = 15; + auto manifest = WriteDeleteManifest( + 2, kSnapshotId, {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, delete_file)}, + unpartitioned_spec_); + auto table = MakePositionDeletesTable( + MakeTable(2, unpartitioned_spec_, MakeSnapshot(2, kSnapshotId, 1, {manifest}))); + + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan()); + auto batch = Import(std::move(array), *table->schema()); + ASSERT_EQ(batch->num_rows(), 1); + + auto rows = std::static_pointer_cast<::arrow::StructArray>(batch->column(2)); + ASSERT_FALSE(rows->IsNull(0)); + auto ids = std::static_pointer_cast<::arrow::Int32Array>(rows->field(0)); + auto data = std::static_pointer_cast<::arrow::StringArray>(rows->field(1)); + EXPECT_EQ(ids->Value(0), 42); + EXPECT_EQ(data->GetString(0), "deleted"); +} + +TEST_F(PositionDeletesTableTest, ProjectsSubsetDeletedRowsFromParquet) { + ICEBERG_UNWRAP_OR_FAIL( + auto delete_file, + WritePositionDeletesWithRows("position-deletes-with-subset-rows.parquet", + R"([["data.parquet", 5, ["deleted"]]])", + {schema_->fields()[1]}, true, unpartitioned_spec_, + PartitionValues{})); + constexpr int64_t kSnapshotId = 16; + auto manifest = WriteDeleteManifest( + 2, kSnapshotId, {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, delete_file)}, + unpartitioned_spec_); + auto table = MakePositionDeletesTable( + MakeTable(2, unpartitioned_spec_, MakeSnapshot(2, kSnapshotId, 1, {manifest}))); + + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan()); + auto batch = Import(std::move(array), *table->schema()); + ASSERT_EQ(batch->num_rows(), 1); + + auto rows = std::static_pointer_cast<::arrow::StructArray>(batch->column(2)); + ASSERT_FALSE(rows->IsNull(0)); + auto ids = std::static_pointer_cast<::arrow::Int32Array>(rows->field(0)); + auto data = std::static_pointer_cast<::arrow::StringArray>(rows->field(1)); + EXPECT_TRUE(ids->IsNull(0)); + EXPECT_EQ(data->GetString(0), "deleted"); +} + +TEST_F(PositionDeletesTableTest, DoesNotDefaultOmittedDeletedRowFields) { + schema_ = std::make_shared( + std::vector{ + SchemaField::MakeRequired(1, "id", int32()) + .WithInitialDefault(std::make_shared(Literal::Int(99))), + SchemaField::MakeRequired(2, "data", string()), + }, + 2); + ICEBERG_UNWRAP_OR_FAIL( + auto delete_file, + WritePositionDeletesWithRows("position-deletes-with-defaulted-subset-row.parquet", + R"([["data.parquet", 5, ["deleted"]]])", + {schema_->fields()[1]}, true, unpartitioned_spec_, + PartitionValues{})); + constexpr int64_t kSnapshotId = 17; + auto manifest = WriteDeleteManifest( + 3, kSnapshotId, {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, delete_file)}, + unpartitioned_spec_); + auto table = MakePositionDeletesTable( + MakeTable(3, unpartitioned_spec_, MakeSnapshot(3, kSnapshotId, 1, {manifest}))); + + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan()); + auto batch = Import(std::move(array), *table->schema()); + ASSERT_EQ(batch->num_rows(), 1); + + auto rows = std::static_pointer_cast<::arrow::StructArray>(batch->column(2)); + auto ids = std::static_pointer_cast<::arrow::Int32Array>(rows->field(0)); + auto data = std::static_pointer_cast<::arrow::StringArray>(rows->field(1)); + EXPECT_TRUE(ids->IsNull(0)); + EXPECT_EQ(data->GetString(0), "deleted"); +} + +TEST_F(PositionDeletesTableTest, RejectsNullDeletedRowsFromParquet) { + ICEBERG_UNWRAP_OR_FAIL( + auto delete_file, + WritePositionDeletesWithRows( + "position-deletes-with-null-row.parquet", R"([["data.parquet", 5, null]])", + std::vector(schema_->fields().begin(), schema_->fields().end()), + false, unpartitioned_spec_, PartitionValues{})); + constexpr int64_t kSnapshotId = 19; + auto manifest = WriteDeleteManifest( + 2, kSnapshotId, {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, delete_file)}, + unpartitioned_spec_); + auto table = MakePositionDeletesTable( + MakeTable(2, unpartitioned_spec_, MakeSnapshot(2, kSnapshotId, 1, {manifest}))); + + auto result = table->Scan(); + EXPECT_THAT(result, IsError(ErrorKind::kInvalidArrowData)); + EXPECT_THAT(result, HasErrorMessage("null row values")); +} + +TEST_F(PositionDeletesTableTest, ReassignsPartitionIdsCollidingWithNestedRows) { + schema_ = std::make_shared(std::vector{ + SchemaField::MakeRequired(1, "payload", + std::make_shared(std::vector{ + SchemaField::MakeRequired(1000, "id", int32()), + })), + SchemaField::MakeRequired(2, "data", string()), + }); + ICEBERG_UNWRAP_OR_FAIL( + auto spec, + PartitionSpec::Make(2, {PartitionField(1000, 1000, "id", Transform::Identity())})); + auto shared_spec = std::shared_ptr(std::move(spec)); + auto table = MakePositionDeletesTable(MakeTable(2, shared_spec)); + + EXPECT_THAT(table->schema()->HighestFieldId(), IsOk()); + const auto& fields = table->schema()->fields(); + auto partition_type = std::static_pointer_cast(fields[3].type()); + ASSERT_EQ(partition_type->fields().size(), 1); + EXPECT_NE(partition_type->fields()[0].field_id(), 1000); + EXPECT_NE(partition_type->fields()[0].field_id(), MetadataColumns::kPartitionColumnId); +} + +TEST_F(PositionDeletesTableTest, RejectsIncompatiblePartitionEvolution) { + ICEBERG_UNWRAP_OR_FAIL( + auto old_spec, + PartitionSpec::Make(1, {PartitionField(2, 1000, "data", Transform::Identity())})); + ICEBERG_UNWRAP_OR_FAIL( + auto new_spec, + PartitionSpec::Make(2, {PartitionField(2, 1000, "data", Transform::Bucket(16))})); + auto shared_old_spec = std::shared_ptr(std::move(old_spec)); + auto shared_new_spec = std::shared_ptr(std::move(new_spec)); + auto source = + MakeTableWithSpecs(2, {shared_old_spec, shared_new_spec}, shared_new_spec); + + auto table = MetadataTable::Make(source); + EXPECT_THAT(table, IsError(ErrorKind::kInvalidSchema)); + EXPECT_THAT(table, HasErrorMessage("incompatible transforms")); +} + +TEST_F(PositionDeletesTableTest, InspectsLogicalPartitionTypes) { + struct Case { + const char* name; + Literal value; + const char* expected; + }; + for (const auto& test : + std::vector{{"decimal", Literal::Decimal(-1234, 10, 2), "-12.34"}, + {"date", Literal::Date(1), "1970-01-02"}}) { + SCOPED_TRACE(test.name); + schema_ = std::make_shared( + std::vector{SchemaField::MakeRequired(1, "key", test.value.type())}); + ICEBERG_UNWRAP_OR_FAIL( + auto spec, + PartitionSpec::Make(1, {PartitionField(1, 1000, "key", Transform::Identity())})); + auto shared_spec = std::shared_ptr(std::move(spec)); + auto files = + WriteDeletionVectors(std::string(test.name) + ".puffin", {{"data.parquet", {2}}}, + shared_spec, PartitionValues(test.value)); + constexpr int64_t kSnapshotId = 19; + auto manifest = WriteDeleteManifest( + 3, kSnapshotId, + {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, files.front())}, shared_spec); + auto table = MakePositionDeletesTable( + MakeTable(3, shared_spec, MakeSnapshot(3, kSnapshotId, 1, {manifest}))); + ICEBERG_UNWRAP_OR_FAIL(auto stream, table->Scan()); + auto batch = Import(std::move(stream), *table->schema()); + ASSERT_EQ(batch->num_rows(), 1); + auto partitions = std::static_pointer_cast<::arrow::StructArray>(batch->column(3)); + auto value = partitions->field(0)->GetScalar(0).ValueOrDie(); + EXPECT_EQ(value->ToString(), test.expected); + } +} + +TEST_F(PositionDeletesTableTest, PromotesHistoricalPartitionValues) { + ICEBERG_UNWRAP_OR_FAIL( + auto spec, + PartitionSpec::Make(1, {PartitionField(1, 1000, "id", Transform::Identity())})); + auto shared_spec = std::shared_ptr(std::move(spec)); + auto historical_schema = schema_; + auto delete_file = WritePositionDeletes("position-deletes-int-partition.parquet", + {{"data.parquet", 5}}, shared_spec, + PartitionValues(Literal::Int(11))); + constexpr int64_t kSnapshotId = 18; + auto manifest = WriteDeleteManifest( + 2, kSnapshotId, {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, delete_file)}, + shared_spec); + + schema_ = std::make_shared( + std::vector{ + SchemaField::MakeRequired(1, "id", int64()), + SchemaField::MakeRequired(2, "data", string()), + }, + 2); + auto source = MakeTable(2, shared_spec, MakeSnapshot(2, kSnapshotId, 1, {manifest})); + source->metadata()->schemas.insert(source->metadata()->schemas.begin(), + historical_schema); + + auto table = MakePositionDeletesTable(source); + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan()); + auto batch = Import(std::move(array), *table->schema()); + ASSERT_EQ(batch->num_rows(), 1); + auto partitions = std::static_pointer_cast<::arrow::StructArray>(batch->column(3)); + auto ids = std::static_pointer_cast<::arrow::Int64Array>(partitions->field(0)); + EXPECT_EQ(ids->Value(0), 11); +} + +TEST_F(PositionDeletesTableTest, OmitsDroppedHistoricalPartitionSources) { + auto historical_schema = std::make_shared( + std::vector{ + SchemaField::MakeRequired(1, "id", int32()), + SchemaField::MakeRequired(2, "data", string()), + SchemaField::MakeOptional(3, "dropped", string()), + }, + 1); + ICEBERG_UNWRAP_OR_FAIL( + auto old_spec, PartitionSpec::Make( + 1, {PartitionField(3, 1000, "dropped", Transform::Identity())})); + auto shared_old_spec = std::shared_ptr(std::move(old_spec)); + auto source = + MakeTableWithSpecs(2, {shared_old_spec, unpartitioned_spec_}, unpartitioned_spec_); + source->metadata()->schemas.insert(source->metadata()->schemas.begin(), + historical_schema); + + auto table = MakePositionDeletesTable(source); + std::vector names; + for (const auto& field : table->schema()->fields()) { + names.emplace_back(field.name()); + } + EXPECT_THAT(names, + ElementsAre("file_path", "pos", "row", "spec_id", "delete_file_path")); +} + +TEST_F(PositionDeletesTableTest, ProjectsEmptyScan) { + auto table = MakePositionDeletesTable(MakeTable(3, unpartitioned_spec_)); + ICEBERG_UNWRAP_OR_FAIL(auto projected, table->schema()->Select(std::vector{ + "pos", "delete_file_path"})); + + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan(*projected)); + auto batch = Import(std::move(array), *projected); + EXPECT_EQ(batch->num_rows(), 0); + ASSERT_EQ(batch->num_columns(), 2); + EXPECT_EQ(batch->schema()->field(0)->name(), "pos"); + EXPECT_EQ(batch->schema()->field(1)->name(), "delete_file_path"); +} + +TEST_F(PositionDeletesTableTest, RejectsNestedRowProjection) { + auto table = MakePositionDeletesTable(MakeTable(3, unpartitioned_spec_)); + ICEBERG_UNWRAP_OR_FAIL(auto projected, + table->schema()->Select(std::vector{"row.id"})); + + auto result = table->Scan(*projected); + EXPECT_THAT(result, IsError(ErrorKind::kInvalidArgument)); + EXPECT_THAT(result, HasErrorMessage("only supports complete top-level fields")); +} + +TEST_F(PositionDeletesTableTest, ExpandsUpgradedMixedTable) { + auto parquet_file = WritePositionDeletes( + "old-position-deletes.parquet", {{"old-data.parquet", 3}, {"old-data.parquet", 9}}, + unpartitioned_spec_, PartitionValues{}); + auto dv_files = + WriteDeletionVectors("new-deletes.puffin", {{"new-data.parquet", {4, 12}}}, + unpartitioned_spec_, PartitionValues{}); + + constexpr int64_t kSnapshotId = 20; + auto old_manifest = WriteDeleteManifest( + 2, kSnapshotId, + {MakeEntry(ManifestStatus::kExisting, kSnapshotId, 1, parquet_file)}, + unpartitioned_spec_); + auto new_manifest = WriteDeleteManifest( + 3, kSnapshotId, + {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 2, dv_files.front())}, + unpartitioned_spec_); + auto table = MakePositionDeletesTable( + MakeTable(3, unpartitioned_spec_, + MakeSnapshot(3, kSnapshotId, 2, {old_manifest, new_manifest}))); + + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan()); + auto batch = Import(std::move(array), *table->schema()); + ASSERT_EQ(batch->num_rows(), 4); + + auto data_paths = std::static_pointer_cast<::arrow::StringArray>(batch->column(0)); + auto positions = std::static_pointer_cast<::arrow::Int64Array>(batch->column(1)); + auto rows = std::static_pointer_cast<::arrow::StructArray>(batch->column(2)); + auto delete_paths = std::static_pointer_cast<::arrow::StringArray>(batch->column(4)); + auto offsets = std::static_pointer_cast<::arrow::Int64Array>(batch->column(5)); + auto sizes = std::static_pointer_cast<::arrow::Int64Array>(batch->column(6)); + + EXPECT_EQ(rows->null_count(), 4); + std::vector actual_data_paths{ + data_paths->GetString(0), data_paths->GetString(1), data_paths->GetString(2), + data_paths->GetString(3)}; + std::vector actual_positions{positions->Value(0), positions->Value(1), + positions->Value(2), positions->Value(3)}; + std::vector actual_delete_paths{delete_paths->GetString(0), + delete_paths->GetString(2)}; + EXPECT_THAT(actual_data_paths, ElementsAre("old-data.parquet", "old-data.parquet", + "new-data.parquet", "new-data.parquet")); + EXPECT_THAT(actual_positions, ElementsAre(3, 9, 4, 12)); + EXPECT_THAT(actual_delete_paths, + ElementsAre("old-position-deletes.parquet", "new-deletes.puffin")); + EXPECT_TRUE(offsets->IsNull(0)); + EXPECT_TRUE(sizes->IsNull(0)); + EXPECT_FALSE(offsets->IsNull(2)); + EXPECT_FALSE(sizes->IsNull(2)); + EXPECT_EQ(offsets->Value(2), *dv_files.front()->content_offset); + EXPECT_EQ(sizes->Value(2), *dv_files.front()->content_size_in_bytes); +} + +TEST_F(PositionDeletesTableTest, RequiresRecreationAfterSourceLayoutChange) { + auto source = MakeTable(2, unpartitioned_spec_); + auto table = MakePositionDeletesTable(source); + auto catalog = std::dynamic_pointer_cast(source->catalog()); + ASSERT_NE(catalog, nullptr); + auto upgraded = MakeTable(3, unpartitioned_spec_); + ICEBERG_UNWRAP_OR_FAIL(upgraded, Table::Make(upgraded->name(), upgraded->metadata(), + "/tmp/table/upgraded.metadata.json", + upgraded->io(), upgraded->catalog())); + EXPECT_CALL(*catalog, LoadTable(::testing::_)).WillOnce(::testing::Return(upgraded)); + ASSERT_THAT(source->Refresh(), IsOk()); + EXPECT_THAT(table->Scan(), HasErrorMessage("recreate the position_deletes table")); + auto recreated = MakePositionDeletesTable(source); + ICEBERG_UNWRAP_OR_FAIL(auto stream, recreated->Scan()); + auto batch = Import(std::move(stream), *recreated->schema()); + EXPECT_EQ(batch->num_rows(), 0); + EXPECT_EQ(batch->num_columns(), 7); +} + +TEST_F(PositionDeletesTableTest, ProjectsAcrossBoundedBatches) { + std::vector positions(MetadataTable::kBatchSize + 7); + std::iota(positions.begin(), positions.end(), 0); + auto files = WriteDeletionVectors("batched.puffin", {{"data.parquet", positions}}, + unpartitioned_spec_, PartitionValues{}); + constexpr int64_t kSnapshotId = 31; + auto manifest = WriteDeleteManifest( + 3, kSnapshotId, {MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, files.front())}, + unpartitioned_spec_); + auto table = MakePositionDeletesTable( + MakeTable(3, unpartitioned_spec_, MakeSnapshot(3, kSnapshotId, 1, {manifest}))); + const auto& fields = table->schema()->fields(); + auto projected = + std::make_unique(std::vector{fields[1], fields[0]}); + ICEBERG_UNWRAP_OR_FAIL(auto stream, table->Scan(*projected)); + auto reader = ::arrow::ImportRecordBatchReader(&stream).ValueOrDie(); + auto first = reader->Next().ValueOrDie(); + ASSERT_NE(first, nullptr); + EXPECT_EQ(first->schema()->field(0)->name(), "pos"); + EXPECT_EQ(first->schema()->field(1)->name(), "file_path"); + EXPECT_EQ(first->num_rows(), MetadataTable::kBatchSize); + auto second = reader->Next().ValueOrDie(); + ASSERT_NE(second, nullptr); + EXPECT_EQ(second->num_rows(), 7); + auto values = std::static_pointer_cast<::arrow::Int64Array>(second->column(0)); + EXPECT_EQ(values->Value(0), MetadataTable::kBatchSize); + EXPECT_EQ(reader->Next().ValueOrDie(), nullptr); +} + +TEST_F(PositionDeletesTableTest, ExpandsMultiplePuffinBlobs) { + auto dv_files = + WriteDeletionVectors("multi-deletes.puffin", + {{"data-a.parquet", {1, 5}}, {"data-b.parquet", {2, 8, 13}}}, + unpartitioned_spec_, PartitionValues{}); + ASSERT_EQ(dv_files.size(), 2); + + constexpr int64_t kSnapshotId = 30; + std::vector entries; + for (const auto& file : dv_files) { + entries.push_back(MakeEntry(ManifestStatus::kAdded, kSnapshotId, 1, file)); + } + auto manifest = + WriteDeleteManifest(3, kSnapshotId, std::move(entries), unpartitioned_spec_); + auto table = MakePositionDeletesTable( + MakeTable(3, unpartitioned_spec_, MakeSnapshot(3, kSnapshotId, 1, {manifest}))); + + ICEBERG_UNWRAP_OR_FAIL(auto array, table->Scan()); + auto batch = Import(std::move(array), *table->schema()); + ASSERT_EQ(batch->num_rows(), 5); + + auto data_paths = std::static_pointer_cast<::arrow::StringArray>(batch->column(0)); + auto positions = std::static_pointer_cast<::arrow::Int64Array>(batch->column(1)); + std::vector> rows; + for (int64_t i = 0; i < batch->num_rows(); ++i) { + rows.emplace_back(data_paths->GetString(i), positions->Value(i)); + } + EXPECT_THAT(rows, ElementsAre(std::tuple{"data-a.parquet", int64_t{1}}, + std::tuple{"data-a.parquet", int64_t{5}}, + std::tuple{"data-b.parquet", int64_t{2}}, + std::tuple{"data-b.parquet", int64_t{8}}, + std::tuple{"data-b.parquet", int64_t{13}})); + EXPECT_NE(dv_files[0]->content_offset, dv_files[1]->content_offset); +} + +INSTANTIATE_TEST_SUITE_P(FormatVersions, PositionDeletesTableTest, + ::testing::Values(2, 3)); + +} // namespace +} // namespace iceberg diff --git a/src/iceberg/type_fwd.h b/src/iceberg/type_fwd.h index f45dd0115..cd3da4e38 100644 --- a/src/iceberg/type_fwd.h +++ b/src/iceberg/type_fwd.h @@ -198,6 +198,7 @@ class ManifestListWriter; class ManifestReader; class ManifestWriter; class PartitionSummary; +class PositionDeletesTable; /// \brief File I/O. struct ReaderOptions;