diff --git a/include/paimon/data/shredding/map_shared_shredding_schema_utils.h b/include/paimon/data/shredding/map_shared_shredding_schema_utils.h index b8470c4d6..bce0b6534 100644 --- a/include/paimon/data/shredding/map_shared_shredding_schema_utils.h +++ b/include/paimon/data/shredding/map_shared_shredding_schema_utils.h @@ -56,6 +56,48 @@ struct PAIMON_EXPORT MapSharedShreddingFieldMeta { } }; +/// Builds a selected-key projection field for a top-level shared-shredding MAP column. +/// +/// The built field is a STRUCT which replaces the MAP field in the read schema. Each child +/// corresponds to one selected key and contains that key's MAP value, or NULL when the key is +/// absent. Children use the selected keys as their names and preserve insertion order. +/// +/// Example: read keys "age" and "score" from MAP column `attributes`: +/// +/// auto builder = MapSharedShreddingAccessBuilder::Create(attributes_field); +/// builder->AddKey("age"); +/// builder->AddKey("score"); +/// auto field = builder->Build(); +/// +/// Use the returned field in `ReadContextBuilder::SetReadSchema`. +class PAIMON_EXPORT MapSharedShreddingAccessBuilder { + public: + /// Creates a builder bound to the original MAP field. + /// + /// The field must be a MAP with STRING keys. Its name, nullability, and value type are + /// retained for the selected-key projection. Ownership of the Arrow C schema resources is + /// transferred to this method. + static Result> Create( + struct ArrowSchema* map_field); + + ~MapSharedShreddingAccessBuilder(); + + /// Adds a selected MAP key. + /// + /// @param key The string MAP key. Keys are returned in insertion order. + Status AddKey(const std::string& key); + + /// Builds a STRUCT projection field which retains the original MAP field's name and + /// nullability. Every selected-key child uses the complete MAP value type and is nullable. + Result> Build() const; + + private: + class Impl; + explicit MapSharedShreddingAccessBuilder(std::unique_ptr&& impl); + + std::unique_ptr impl_; +}; + class PAIMON_EXPORT MapSharedShreddingSchemaUtils { public: MapSharedShreddingSchemaUtils() = delete; diff --git a/include/paimon/read_context.h b/include/paimon/read_context.h index 8fdac7b35..91a1033fc 100644 --- a/include/paimon/read_context.h +++ b/include/paimon/read_context.h @@ -226,6 +226,12 @@ class PAIMON_EXPORT ReadContextBuilder { /// key list, for example: "k1,k2". Only map fields with string key type /// (Arrow utf8) are supported. /// + /// Attaching this metadata to a MAP field only filters the returned MAP after + /// reading. To push down selected keys from a shared-shredding MAP and return + /// them as STRUCT children, build the field with + /// `MapSharedShreddingAccessBuilder`. To read selected paths from a VARIANT + /// field, build the field with `VariantAccessBuilder`. + /// /// Example: /// @code{.cpp} /// auto map_field = arrow::field("m", arrow::map(arrow::utf8(), arrow::int32())); diff --git a/src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp b/src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp index dcceddc71..342c6adb4 100644 --- a/src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp +++ b/src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp @@ -21,6 +21,7 @@ #include #include +#include #include #include @@ -30,18 +31,211 @@ #include "paimon/common/reader/reader_utils.h" #include "paimon/common/utils/arrow/mem_utils.h" #include "paimon/common/utils/arrow/status_utils.h" -#include "paimon/common/utils/string_utils.h" #include "paimon/core/casting/casting_utils.h" +#include "paimon/core/utils/nested_projection_utils.h" namespace paimon { +namespace { + +std::vector> ResolveSelectedKeyIds( + const MapSharedShreddingFieldMeta& meta, const std::vector& selected_keys) { + std::vector> selected_key_ids; + selected_key_ids.reserve(selected_keys.size()); + for (const auto& selected_key : selected_keys) { + auto id_iter = meta.name_to_id.find(selected_key); + if (id_iter != meta.name_to_id.end()) { + selected_key_ids.emplace_back(selected_key, id_iter->second); + } + } + return selected_key_ids; +} + +void CollectPhysicalColumns( + const std::shared_ptr& physical_struct_array, + std::map>* physical_column_name_to_array, + std::shared_ptr* overflow_array) { + const auto& struct_type = physical_struct_array->struct_type(); + for (int32_t i = 0; i < struct_type->num_fields(); ++i) { + const auto& sub_field = struct_type->field(i); + if (sub_field->name() == MapSharedShreddingDefine::kFieldMapping) { + continue; + } + if (sub_field->name() == MapSharedShreddingDefine::kOverflow) { + *overflow_array = arrow::internal::checked_pointer_cast( + physical_struct_array->field(i)); + continue; + } + (*physical_column_name_to_array)[sub_field->name()] = physical_struct_array->field(i); + } +} + +class FullMapReadPlan : public MapFieldReadPlan { + public: + FullMapReadPlan(const std::shared_ptr& logical_field, + const std::shared_ptr& physical_read_field, + std::vector>&& selected_key_ids) + : MapFieldReadPlan(logical_field, physical_read_field), + selected_key_ids_(std::move(selected_key_ids)), + logical_map_type_( + arrow::internal::checked_pointer_cast(logical_field->type())) {} + + Result> Materialize( + const std::shared_ptr& physical_array, + arrow::MemoryPool* arrow_pool) const override; + + private: + std::vector> selected_key_ids_; + std::shared_ptr logical_map_type_; +}; + +class SharedSelectedKeysReadPlan : public MapFieldReadPlan { + public: + struct SelectedKey { + int32_t field_id = -1; + std::vector candidate_columns; + bool may_use_overflow = false; + }; + + SharedSelectedKeysReadPlan(const std::shared_ptr& logical_field, + const std::shared_ptr& physical_read_field, + std::vector&& selected_keys) + : MapFieldReadPlan(logical_field, physical_read_field), + selected_keys_(std::move(selected_keys)) {} + + Result> Materialize( + const std::shared_ptr& physical_array, + arrow::MemoryPool* arrow_pool) const override; + + private: + std::vector selected_keys_; +}; + +class DefaultSelectedKeysReadPlan : public MapFieldReadPlan { + public: + DefaultSelectedKeysReadPlan(const std::shared_ptr& logical_field, + const std::shared_ptr& physical_read_field, + const std::vector& selected_keys) + : MapFieldReadPlan(logical_field, physical_read_field), selected_keys_(selected_keys) {} + + Result> Materialize( + const std::shared_ptr& physical_array, + arrow::MemoryPool* arrow_pool) const override; + + private: + std::vector selected_keys_; +}; + +} // namespace + +Result> MapFieldReadPlanFactory::CreateMapReadPlan( + const std::shared_ptr& logical_map_field, + const MapSharedShreddingFieldMeta& meta) { + if (logical_map_field->type()->id() != arrow::Type::MAP) { + return Status::Invalid(fmt::format("full MAP read plan requires MAP field {}, got {}", + logical_map_field->name(), + logical_map_field->type()->ToString())); + } + auto logical_map_type = + arrow::internal::checked_pointer_cast(logical_map_field->type()); + PAIMON_ASSIGN_OR_RAISE(std::vector selected_keys, + NestedProjectionUtils::GetMapSelectedKeys(logical_map_field)); + if (selected_keys.empty()) { + selected_keys.reserve(meta.name_to_id.size()); + for (const auto& [key_name, _] : meta.name_to_id) { + selected_keys.push_back(key_name); + } + } + std::set selected_physical_column_ids; + bool include_overflow = false; + for (const auto& selected_key : selected_keys) { + auto field_id_iter = meta.name_to_id.find(selected_key); + if (field_id_iter == meta.name_to_id.end()) { + continue; + } + int32_t field_id = field_id_iter->second; + include_overflow = include_overflow || meta.overflow_field_set.count(field_id) > 0; + auto columns_iter = meta.field_to_columns.find(field_id); + if (columns_iter != meta.field_to_columns.end()) { + selected_physical_column_ids.insert(columns_iter->second.begin(), + columns_iter->second.end()); + } + } + std::shared_ptr physical_type = + MapSharedShreddingUtils::BuildSpecificPhysicalStructType( + logical_map_type->item_type(), selected_physical_column_ids, + logical_map_type->item_field()->nullable(), include_overflow); + auto physical_read_field = logical_map_field->WithType(physical_type); + std::unique_ptr read_plan = std::make_unique( + logical_map_field, physical_read_field, ResolveSelectedKeyIds(meta, selected_keys)); + return read_plan; +} + +Result> MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan( + const std::shared_ptr& selected_keys_field, + const MapSharedShreddingFieldMeta& meta) { + PAIMON_ASSIGN_OR_RAISE( + std::vector selected_keys, + NestedProjectionUtils::ValidateMapSharedShreddingAccessField(selected_keys_field)); + auto selected_keys_type = + arrow::internal::checked_pointer_cast(selected_keys_field->type()); + const auto& value_field = selected_keys_type->field(0); + + std::set selected_physical_column_ids; + bool include_overflow = false; + std::vector selected_key_plans; + selected_key_plans.reserve(selected_keys.size()); + for (const auto& selected_key : selected_keys) { + SharedSelectedKeysReadPlan::SelectedKey selected_key_plan; + auto field_id_iter = meta.name_to_id.find(selected_key); + if (field_id_iter != meta.name_to_id.end()) { + selected_key_plan.field_id = field_id_iter->second; + auto columns_iter = meta.field_to_columns.find(selected_key_plan.field_id); + if (columns_iter != meta.field_to_columns.end()) { + selected_key_plan.candidate_columns = columns_iter->second; + selected_physical_column_ids.insert(columns_iter->second.begin(), + columns_iter->second.end()); + } + selected_key_plan.may_use_overflow = + meta.overflow_field_set.count(selected_key_plan.field_id) > 0; + include_overflow = include_overflow || selected_key_plan.may_use_overflow; + } + selected_key_plans.push_back(std::move(selected_key_plan)); + } + std::shared_ptr physical_type = + MapSharedShreddingUtils::BuildSpecificPhysicalStructType( + value_field->type(), selected_physical_column_ids, value_field->nullable(), + include_overflow); + auto physical_read_field = selected_keys_field->WithType(physical_type); + std::unique_ptr read_plan = std::make_unique( + selected_keys_field, physical_read_field, std::move(selected_key_plans)); + return read_plan; +} + +Result> +MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan( + const std::shared_ptr& file_map_field, + const std::shared_ptr& selected_keys_field) { + if (file_map_field->type()->id() != arrow::Type::MAP) { + return Status::Invalid( + fmt::format("selected-key MAP projection {} requires MAP file field, got {}", + selected_keys_field->name(), file_map_field->type()->ToString())); + } + PAIMON_ASSIGN_OR_RAISE( + std::vector selected_keys, + NestedProjectionUtils::ValidateMapSharedShreddingAccessField(selected_keys_field)); + auto physical_read_field = selected_keys_field->WithType(file_map_field->type()); + std::unique_ptr read_plan = std::make_unique( + selected_keys_field, physical_read_field, selected_keys); + return read_plan; +} + MapSharedShreddingFileReader::MapSharedShreddingFileReader( std::unique_ptr&& reader, - std::map&& - shared_shredding_name_to_context, + std::map>&& field_read_plans, const std::shared_ptr& pool) : arrow_pool_(GetArrowPool(pool)), reader_(std::move(reader)), - shared_shredding_name_to_context_(std::move(shared_shredding_name_to_context)) {} + field_read_plans_(std::move(field_read_plans)) {} Result> MapSharedShreddingFileReader::GetFileSchema() const { PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<::ArrowSchema> physical_schema, @@ -68,7 +262,8 @@ Result> MapSharedShreddingFileReader::GetFileSche Result> MapSharedShreddingFileReader::ToLogicalMapField( const std::shared_ptr& physical_field) { - auto physical_type = std::dynamic_pointer_cast(physical_field->type()); + auto physical_type = + arrow::internal::checked_pointer_cast(physical_field->type()); if (!physical_type) { return Status::Invalid(fmt::format("shared-shredding field {} is not a physical struct", physical_field->name())); @@ -103,62 +298,24 @@ Status MapSharedShreddingFileReader::SetReadSchema( } PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr logical_read_schema, arrow::ImportSchema(read_schema)); - std::vector shared_shredding_names; - for (const auto& field : logical_read_schema->fields()) { - if (shared_shredding_name_to_context_.find(field->name()) != - shared_shredding_name_to_context_.end()) { - shared_shredding_names.push_back(field->name()); + bool converted = false; + arrow::FieldVector physical_read_fields = logical_read_schema->fields(); + for (size_t i = 0; i < logical_read_schema->fields().size(); ++i) { + const auto& field = logical_read_schema->field(i); + auto plan_iter = field_read_plans_.find(field->name()); + if (plan_iter != field_read_plans_.end()) { + physical_read_fields[i] = plan_iter->second->PhysicalReadField(); + converted = true; } } - if (shared_shredding_names.empty()) { - // suppose not fall into MapSharedShreddingFileReader - return Status::Invalid("do not exist shared shredding columns in read schema"); + if (!converted) { + return Status::Invalid("suppose not fall into MapSharedShreddingFileReader"); } - arrow::FieldVector resolved_fields = logical_read_schema->fields(); - for (const auto& name : shared_shredding_names) { - const auto& field = logical_read_schema->GetFieldByName(name); - if (!field) { - return Status::Invalid( - fmt::format("cannot find shared-shredding field {} in read schema", name)); - } - auto context_iter = shared_shredding_name_to_context_.find(field->name()); - if (context_iter == shared_shredding_name_to_context_.end()) { - return Status::Invalid( - fmt::format("cannot find shared-shredding metadata for field {}", field->name())); - } - std::set selected_physical_column_ids; - bool include_overflow = false; - for (const auto& selected_key : context_iter->second.selected_keys) { - // check if selected_key in file - auto name_iter = context_iter->second.meta.name_to_id.find(selected_key); - if (name_iter == context_iter->second.meta.name_to_id.end()) { - continue; - } - // check if selected_key in overflow_field - PAIMON_ASSIGN_OR_RAISE( - bool is_overflow_field, - MapSharedShreddingUtils::IsOverflowField(context_iter->second.meta, selected_key)); - include_overflow = include_overflow || is_overflow_field; - // check if selected_key in field_to_columns - auto column_iter = context_iter->second.meta.field_to_columns.find(name_iter->second); - if (column_iter == context_iter->second.meta.field_to_columns.end()) { - continue; - } - const std::vector& physical_column_ids = column_iter->second; - selected_physical_column_ids.insert(physical_column_ids.begin(), - physical_column_ids.end()); - } - std::shared_ptr resolved_type = - MapSharedShreddingUtils::BuildSpecificPhysicalStructType( - context_iter->second.map_type->item_type(), selected_physical_column_ids, - context_iter->second.map_type->item_field()->nullable(), include_overflow); - resolved_fields[logical_read_schema->GetFieldIndex(name)] = - arrow::field(field->name(), resolved_type, field->nullable()); - } - auto resolved_schema = arrow::schema(std::move(resolved_fields)); - std::unique_ptr c_resolved_schema = std::make_unique(); - PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*resolved_schema, c_resolved_schema.get())); - return reader_->SetReadSchema(c_resolved_schema.get(), predicate, selection_bitmap); + auto physical_read_schema = arrow::schema(std::move(physical_read_fields)); + std::unique_ptr c_physical_read_schema = std::make_unique(); + PAIMON_RETURN_NOT_OK_FROM_ARROW( + arrow::ExportSchema(*physical_read_schema, c_physical_read_schema.get())); + return reader_->SetReadSchema(c_physical_read_schema.get(), predicate, selection_bitmap); } Result MapSharedShreddingFileReader::NextBatch() { @@ -177,7 +334,7 @@ Result MapSharedShreddingFileReader::NextBatch auto& [c_array, c_schema] = batch; PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr arrow_array, arrow::ImportArray(c_array.get(), c_schema.get())); - auto struct_array = std::dynamic_pointer_cast(arrow_array); + auto struct_array = arrow::internal::checked_pointer_cast(arrow_array); if (!struct_array) { return Status::Invalid("cannot cast batch to StructArray in MapSharedShreddingFileReader"); } @@ -186,21 +343,14 @@ Result MapSharedShreddingFileReader::NextBatch arrow::FieldVector resolved_fields = struct_array->struct_type()->fields(); for (int32_t field_idx = 0; field_idx < struct_array->num_fields(); ++field_idx) { const auto& physical_field = struct_array->struct_type()->field(field_idx); - auto iter = shared_shredding_name_to_context_.find(physical_field->name()); - if (iter == shared_shredding_name_to_context_.end()) { + auto plan_iter = field_read_plans_.find(physical_field->name()); + if (plan_iter == field_read_plans_.end()) { continue; } - auto physical_struct_array = - std::dynamic_pointer_cast(struct_array->field(field_idx)); - if (!physical_struct_array) { - return Status::Invalid(fmt::format( - "cannot cast physical shredding field {} to StructArray", physical_field->name())); - } - PAIMON_ASSIGN_OR_RAISE(std::shared_ptr logical_map_array, - RebuildLogicalMapArray(physical_field, physical_struct_array)); - resolved_arrays[field_idx] = logical_map_array; - resolved_fields[field_idx] = arrow::field(physical_field->name(), logical_map_array->type(), - physical_field->nullable()); + PAIMON_ASSIGN_OR_RAISE( + resolved_arrays[field_idx], + plan_iter->second->Materialize(struct_array->field(field_idx), arrow_pool_.get())); + resolved_fields[field_idx] = plan_iter->second->LogicalField(); } PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr new_struct_array, arrow::StructArray::Make(resolved_arrays, resolved_fields)); @@ -212,32 +362,28 @@ Result MapSharedShreddingFileReader::NextBatch return batch_with_bitmap; } -Result> MapSharedShreddingFileReader::RebuildLogicalMapArray( - const std::shared_ptr& physical_field, - const std::shared_ptr& physical_struct_array) const { - std::string shredding_field_name = physical_field->name(); - auto iter = shared_shredding_name_to_context_.find(shredding_field_name); - if (iter == shared_shredding_name_to_context_.end()) { - return Status::Invalid( - fmt::format("cannot find shared-shredding context for field {}", shredding_field_name)); +Result> FullMapReadPlan::Materialize( + const std::shared_ptr& physical_array, arrow::MemoryPool* arrow_pool) const { + auto physical_struct_array = + arrow::internal::checked_pointer_cast(physical_array); + if (!physical_struct_array) { + return Status::Invalid(fmt::format("cannot cast physical shredding field {} to StructArray", + LogicalField()->name())); } - const MapSharedShreddingFieldMeta& meta = iter->second.meta; - const std::vector& selected_keys = iter->second.selected_keys; - const auto& map_type = iter->second.map_type; + const std::string& shredding_field_name = LogicalField()->name(); - auto field_mapping_array = std::dynamic_pointer_cast( + auto field_mapping_array = arrow::internal::checked_pointer_cast( physical_struct_array->GetFieldByName(MapSharedShreddingDefine::kFieldMapping)); if (!field_mapping_array) { return Status::Invalid( fmt::format("cannot find __field_mapping for field {}", shredding_field_name)); } auto field_mapping_values = - std::dynamic_pointer_cast(field_mapping_array->values()); + arrow::internal::checked_pointer_cast(field_mapping_array->values()); if (!field_mapping_values) { return Status::Invalid("__field_mapping values is not an Int32Array"); } - auto selected_key_ids = ResolveSelectedKeyIds(meta, selected_keys); std::map> physical_column_name_to_array; std::shared_ptr overflow_array; CollectPhysicalColumns(physical_struct_array, &physical_column_name_to_array, &overflow_array); @@ -245,8 +391,8 @@ Result> MapSharedShreddingFileReader::RebuildLogic if (physical_column_array->type_id() == arrow::Type::DICTIONARY) { PAIMON_ASSIGN_OR_RAISE( physical_column_array, - CastingUtils::Cast(physical_column_array, map_type->item_type(), - arrow::compute::CastOptions::Safe(), arrow_pool_.get())); + CastingUtils::Cast(physical_column_array, logical_map_type_->item_type(), + arrow::compute::CastOptions::Safe(), arrow_pool)); } } @@ -262,13 +408,13 @@ Result> MapSharedShreddingFileReader::RebuildLogic if (overflow_items->type_id() == arrow::Type::DICTIONARY) { PAIMON_ASSIGN_OR_RAISE( overflow_items, - CastingUtils::Cast(overflow_items, map_type->item_type(), - arrow::compute::CastOptions::Safe(), arrow_pool_.get())); + CastingUtils::Cast(overflow_items, logical_map_type_->item_type(), + arrow::compute::CastOptions::Safe(), arrow_pool)); } } PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::unique_ptr map_builder_base, - arrow::MakeBuilder(map_type, arrow_pool_.get())); + arrow::MakeBuilder(logical_map_type_, arrow_pool)); auto* map_builder = dynamic_cast(map_builder_base.get()); if (!map_builder) { return Status::Invalid( @@ -286,7 +432,7 @@ Result> MapSharedShreddingFileReader::RebuildLogic } int64_t row_count = physical_struct_array->length(); - int64_t max_item_count = row_count * static_cast(selected_key_ids.size()); + int64_t max_item_count = row_count * static_cast(selected_key_ids_.size()); PAIMON_RETURN_NOT_OK_FROM_ARROW(map_builder->Reserve(row_count)); PAIMON_RETURN_NOT_OK_FROM_ARROW(key_builder->Reserve(max_item_count)); PAIMON_RETURN_NOT_OK_FROM_ARROW(item_builder->Reserve(max_item_count)); @@ -306,7 +452,7 @@ Result> MapSharedShreddingFileReader::RebuildLogic int32_t mapping_offset = field_mapping_array->value_offset(row); int32_t mapping_length = field_mapping_array->value_length(row); // follow the sequence in paimon.map.selected-keys - for (const auto& [selected_key, selected_field_id] : selected_key_ids) { + for (const auto& [selected_key, selected_field_id] : selected_key_ids_) { bool found = false; for (int32_t pos = 0; pos < mapping_length; ++pos) { int32_t mapping_index = mapping_offset + pos; @@ -353,37 +499,197 @@ Result> MapSharedShreddingFileReader::RebuildLogic return map_array; } -std::vector> MapSharedShreddingFileReader::ResolveSelectedKeyIds( - const MapSharedShreddingFieldMeta& meta, const std::vector& selected_keys) { - std::vector> selected_key_ids; - selected_key_ids.reserve(selected_keys.size()); - for (const auto& selected_key : selected_keys) { - auto id_iter = meta.name_to_id.find(selected_key); - if (id_iter == meta.name_to_id.end()) { +Result> SharedSelectedKeysReadPlan::Materialize( + const std::shared_ptr& physical_array, arrow::MemoryPool* arrow_pool) const { + auto physical_struct_array = + arrow::internal::checked_pointer_cast(physical_array); + if (!physical_struct_array) { + return Status::Invalid(fmt::format("cannot cast physical shredding field {} to StructArray", + LogicalField()->name())); + } + auto selected_keys_type = + arrow::internal::checked_pointer_cast(LogicalField()->type()); + + auto field_mapping_array = arrow::internal::checked_pointer_cast( + physical_struct_array->GetFieldByName(MapSharedShreddingDefine::kFieldMapping)); + if (!field_mapping_array) { + return Status::Invalid( + fmt::format("cannot find __field_mapping for field {}", LogicalField()->name())); + } + auto field_mapping_values = + arrow::internal::checked_pointer_cast(field_mapping_array->values()); + if (!field_mapping_values) { + return Status::Invalid("__field_mapping values is not an Int32Array"); + } + + std::shared_ptr value_type = selected_keys_type->field(0)->type(); + std::map> physical_column_name_to_array; + std::shared_ptr overflow_array; + CollectPhysicalColumns(physical_struct_array, &physical_column_name_to_array, &overflow_array); + for (auto& [_, physical_column_array] : physical_column_name_to_array) { + if (physical_column_array->type_id() == arrow::Type::DICTIONARY) { + PAIMON_ASSIGN_OR_RAISE( + physical_column_array, + CastingUtils::Cast(physical_column_array, value_type, + arrow::compute::CastOptions::Safe(), arrow_pool)); + } + } + + std::shared_ptr overflow_keys; + std::shared_ptr overflow_items; + if (overflow_array) { + overflow_keys = + arrow::internal::checked_pointer_cast(overflow_array->keys()); + overflow_items = overflow_array->items(); + if (!overflow_keys || !overflow_items) { + return Status::Invalid("__overflow map has invalid key or item array"); + } + if (overflow_items->type_id() == arrow::Type::DICTIONARY) { + PAIMON_ASSIGN_OR_RAISE( + overflow_items, + CastingUtils::Cast(overflow_items, value_type, arrow::compute::CastOptions::Safe(), + arrow_pool)); + } + } + + std::unique_ptr access_builder_base; + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(access_builder_base, + arrow::MakeBuilder(LogicalField()->type(), arrow_pool)); + auto* access_builder = dynamic_cast(access_builder_base.get()); + if (!access_builder) { + return Status::Invalid( + fmt::format("selected-key MAP field {} is not a STRUCT", LogicalField()->name())); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Reserve(physical_struct_array->length())); + + for (int64_t row = 0; row < physical_struct_array->length(); ++row) { + if (physical_struct_array->IsNull(row)) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->AppendNull()); continue; } - selected_key_ids.emplace_back(selected_key, id_iter->second); + if (field_mapping_array->IsNull(row)) { + return Status::Invalid(fmt::format( + "__field_mapping cannot be null in non-null shared-shredding row for field {}", + LogicalField()->name())); + } + int32_t mapping_offset = field_mapping_array->value_offset(row); + + PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Append()); + for (int32_t key_index = 0; key_index < selected_keys_type->num_fields(); ++key_index) { + arrow::ArrayBuilder* value_builder = access_builder->field_builder(key_index); + const SelectedKey& selected_key = selected_keys_[key_index]; + bool appended = false; + if (selected_key.field_id >= 0) { + for (int32_t physical_column_id : selected_key.candidate_columns) { + int32_t mapping_index = mapping_offset + physical_column_id; + if (field_mapping_values->IsNull(mapping_index)) { + return Status::Invalid("__field_mapping element cannot be null"); + } + if (field_mapping_values->Value(mapping_index) != selected_key.field_id) { + continue; + } + std::string physical_column_name = + MapSharedShreddingDefine::PhysicalColumnName(physical_column_id); + auto physical_column_iter = + physical_column_name_to_array.find(physical_column_name); + if (physical_column_iter == physical_column_name_to_array.end()) { + return Status::Invalid( + fmt::format("cannot find selected physical column {} for field {}", + physical_column_name, LogicalField()->name())); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->AppendArraySlice( + *physical_column_iter->second->data(), row, 1)); + appended = true; + break; + } + } + + if (!appended && selected_key.may_use_overflow && overflow_array && + !overflow_array->IsNull(row)) { + int32_t overflow_offset = overflow_array->value_offset(row); + int32_t overflow_length = overflow_array->value_length(row); + for (int32_t pos = 0; pos < overflow_length; ++pos) { + int32_t overflow_index = overflow_offset + pos; + if (!overflow_keys->IsNull(overflow_index) && + overflow_keys->Value(overflow_index) == selected_key.field_id) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->AppendArraySlice( + *overflow_items->data(), overflow_index, 1)); + appended = true; + break; + } + } + } + if (!appended) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->AppendNull()); + } + } } - return selected_key_ids; + std::shared_ptr result; + PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Finish(&result)); + return result; } -void MapSharedShreddingFileReader::CollectPhysicalColumns( - const std::shared_ptr& physical_struct_array, - std::map>* physical_column_name_to_array, - std::shared_ptr* overflow_array) { - const auto& struct_type = physical_struct_array->struct_type(); - for (int32_t i = 0; i < struct_type->num_fields(); ++i) { - const auto& sub_field = struct_type->field(i); - if (sub_field->name() == MapSharedShreddingDefine::kFieldMapping) { +Result> DefaultSelectedKeysReadPlan::Materialize( + const std::shared_ptr& physical_array, arrow::MemoryPool* arrow_pool) const { + auto map_array = arrow::internal::checked_pointer_cast(physical_array); + if (!map_array) { + return Status::Invalid( + fmt::format("cannot cast default-layout selected-key field {} to " + "MapArray", + LogicalField()->name())); + } + auto selected_keys_type = + arrow::internal::checked_pointer_cast(LogicalField()->type()); + auto physical_map_type = + arrow::internal::checked_pointer_cast(PhysicalReadField()->type()); + + std::shared_ptr items = map_array->items(); + if (items->type_id() == arrow::Type::DICTIONARY) { + PAIMON_ASSIGN_OR_RAISE(items, + CastingUtils::Cast(items, physical_map_type->item_type(), + arrow::compute::CastOptions::Safe(), arrow_pool)); + } + std::shared_ptr keys = map_array->keys(); + std::unique_ptr access_builder_base; + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(access_builder_base, + arrow::MakeBuilder(LogicalField()->type(), arrow_pool)); + auto* access_builder = dynamic_cast(access_builder_base.get()); + if (!access_builder) { + return Status::Invalid( + fmt::format("selected-key MAP field {} is not a STRUCT", LogicalField()->name())); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Reserve(map_array->length())); + + for (int64_t row = 0; row < map_array->length(); ++row) { + if (map_array->IsNull(row)) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->AppendNull()); continue; } - if (sub_field->name() == MapSharedShreddingDefine::kOverflow) { - *overflow_array = arrow::internal::checked_pointer_cast( - physical_struct_array->field(i)); - continue; + PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Append()); + int64_t begin = map_array->value_offset(row); + int64_t end = map_array->value_offset(row + 1); + for (int32_t key_index = 0; key_index < selected_keys_type->num_fields(); ++key_index) { + arrow::ArrayBuilder* value_builder = access_builder->field_builder(key_index); + bool appended = false; + for (int64_t entry = begin; entry < end; ++entry) { + PAIMON_ASSIGN_OR_RAISE(std::string_view key, + NestedProjectionUtils::GetMapKeyViewAt(keys, entry)); + if (key != selected_keys_[key_index]) { + continue; + } + PAIMON_RETURN_NOT_OK_FROM_ARROW( + value_builder->AppendArraySlice(*items->data(), entry, 1)); + appended = true; + break; + } + if (!appended) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->AppendNull()); + } } - (*physical_column_name_to_array)[sub_field->name()] = physical_struct_array->field(i); } + std::shared_ptr result; + PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Finish(&result)); + return result; } std::shared_ptr MapSharedShreddingFileReader::GetReaderMetrics() const { diff --git a/src/paimon/common/data/shredding/map_shared_shredding_file_reader.h b/src/paimon/common/data/shredding/map_shared_shredding_file_reader.h index ba608cf09..744d09fd1 100644 --- a/src/paimon/common/data/shredding/map_shared_shredding_file_reader.h +++ b/src/paimon/common/data/shredding/map_shared_shredding_file_reader.h @@ -33,21 +33,51 @@ namespace paimon { -class MapSharedShreddingFileReader : public FileBatchReader { +class MapFieldReadPlan { + public: + virtual ~MapFieldReadPlan() = default; + + MapFieldReadPlan(const std::shared_ptr& logical_field, + const std::shared_ptr& physical_read_field) + : logical_field_(logical_field), physical_read_field_(physical_read_field) {} + + const std::shared_ptr& LogicalField() const { + return logical_field_; + } + + const std::shared_ptr& PhysicalReadField() const { + return physical_read_field_; + } + + virtual Result> Materialize( + const std::shared_ptr& physical_array, + arrow::MemoryPool* arrow_pool) const = 0; + + private: + std::shared_ptr logical_field_; + std::shared_ptr physical_read_field_; +}; + +class MapFieldReadPlanFactory { public: - struct SharedShreddingContext { - SharedShreddingContext(const MapSharedShreddingFieldMeta& _meta, - const std::vector& _selected_keys, - const std::shared_ptr& _map_type) - : meta(_meta), selected_keys(_selected_keys), map_type(_map_type) {} - MapSharedShreddingFieldMeta meta; - std::vector selected_keys; - std::shared_ptr map_type; - }; + static Result> CreateMapReadPlan( + const std::shared_ptr& logical_map_field, + const MapSharedShreddingFieldMeta& meta); + + static Result> CreateSharedSelectedKeysReadPlan( + const std::shared_ptr& selected_keys_field, + const MapSharedShreddingFieldMeta& meta); + static Result> CreateDefaultSelectedKeysReadPlan( + const std::shared_ptr& file_map_field, + const std::shared_ptr& selected_keys_field); +}; + +class MapSharedShreddingFileReader : public FileBatchReader { + public: MapSharedShreddingFileReader( std::unique_ptr&& reader, - std::map&& shared_shredding_name_to_context, + std::map>&& field_read_plans, const std::shared_ptr& pool); Result> GetFileSchema() const override; @@ -70,25 +100,13 @@ class MapSharedShreddingFileReader : public FileBatchReader { bool SupportPreciseBitmapSelection() const override; private: - Result> RebuildLogicalMapArray( - const std::shared_ptr& physical_field, - const std::shared_ptr& physical_struct_array) const; - - static std::vector> ResolveSelectedKeyIds( - const MapSharedShreddingFieldMeta& meta, const std::vector& selected_keys); - - static void CollectPhysicalColumns( - const std::shared_ptr& physical_struct_array, - std::map>* physical_column_name_to_array, - std::shared_ptr* overflow_array); - static Result> ToLogicalMapField( const std::shared_ptr& physical_field); private: std::shared_ptr arrow_pool_; std::unique_ptr reader_; - std::map shared_shredding_name_to_context_; + std::map> field_read_plans_; }; } // namespace paimon diff --git a/src/paimon/common/data/shredding/map_shared_shredding_file_reader_test.cpp b/src/paimon/common/data/shredding/map_shared_shredding_file_reader_test.cpp index c8ab3d838..408469e3b 100644 --- a/src/paimon/common/data/shredding/map_shared_shredding_file_reader_test.cpp +++ b/src/paimon/common/data/shredding/map_shared_shredding_file_reader_test.cpp @@ -104,8 +104,7 @@ class MapSharedShreddingFileReaderTest : public ::testing::Test { const std::optional& selected_keys_str = std::nullopt) const { EXPECT_OK_AND_ASSIGN(auto c_file_schema, reader->GetFileSchema()); auto file_schema = arrow::ImportSchema(c_file_schema.get()).ValueOrDie(); - std::map - shared_shredding_name_to_context; + std::map> field_read_plans; for (const auto& field : file_schema->fields()) { auto metadata = std::const_pointer_cast(field->metadata()); if (!MapSharedShreddingUtils::HasShreddingMetadata(metadata)) { @@ -125,22 +124,17 @@ class MapSharedShreddingFileReaderTest : public ::testing::Test { EXPECT_TRUE(item_field); auto map_type = arrow::internal::checked_pointer_cast(arrow::map( arrow::utf8(), arrow::field("value", item_field->type(), item_field->nullable()))); - std::vector selected_keys; + std::shared_ptr logical_map_field = field->WithType(map_type); if (selected_keys_str.has_value()) { - selected_keys = StringUtils::Split(selected_keys_str.value(), ",", - /*ignore_empty=*/false); - } else { - selected_keys.reserve(meta.name_to_id.size()); - for (const auto& [key_name, _] : meta.name_to_id) { - selected_keys.push_back(key_name); - } + logical_map_field = logical_map_field->WithMetadata(arrow::KeyValueMetadata::Make( + {DataField::MAP_SELECTED_KEYS}, {selected_keys_str.value()})); } - shared_shredding_name_to_context.emplace( - field->name(), MapSharedShreddingFileReader::SharedShreddingContext( - meta, selected_keys, map_type)); + EXPECT_OK_AND_ASSIGN(auto field_read_plan, MapFieldReadPlanFactory::CreateMapReadPlan( + logical_map_field, meta)); + field_read_plans.emplace(field->name(), std::move(field_read_plan)); } - return std::make_unique( - std::move(reader), std::move(shared_shredding_name_to_context), pool_); + return std::make_unique(std::move(reader), + std::move(field_read_plans), pool_); } Result> CreateReader( @@ -299,6 +293,118 @@ TEST_F(MapSharedShreddingFileReaderTest, TestAllExistSelectedKeysWithOverflow) { AssertChunkedArrayEquals(expected, actual); } +TEST_F(MapSharedShreddingFileReaderTest, TestSelectedKeysStructProjection) { + ASSERT_OK_AND_ASSIGN(auto physical_schema, PhysicalSchemaWithMetadata()); + ASSERT_OK_AND_ASSIGN(auto physical_array, PhysicalArray()); + auto mock_reader = std::make_unique( + physical_array, arrow::struct_(physical_schema->fields()), /*read_batch_size=*/10); + mock_reader->EnableRandomizeBatchSize(false); + + auto selected_type = + arrow::struct_({arrow::field("a", arrow::int64()), arrow::field("c", arrow::int64()), + arrow::field("missing", arrow::int64())}); + auto selected_field = arrow::field( + "tags", selected_type, /*nullable=*/true, + arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,c,missing"})); + ASSERT_OK_AND_ASSIGN( + auto field_read_plan, + MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(selected_field, TagsMeta())); + std::map> contexts; + contexts.emplace("tags", std::move(field_read_plan)); + auto reader = std::make_unique(std::move(mock_reader), + std::move(contexts), pool_); + + auto read_schema = + ExportSchema(arrow::schema({arrow::field("id", arrow::int32()), selected_field})); + ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr, + /*selection_bitmap=*/std::nullopt)); + ASSERT_OK_AND_ASSIGN(auto actual, ReadResultCollector::CollectResult(reader.get())); + + auto expected_type = arrow::struct_({arrow::field("id", arrow::int32()), selected_field}); + std::shared_ptr expected; + ASSERT_TRUE(arrow::ipc::internal::json::ChunkedArrayFromJSON(expected_type, {R"([ + [1, [10, null, null]], + [2, [40, 30, null]], + [3, null], + [4, [80, null, null]] + ])"}, + &expected) + .ok()); + AssertChunkedArrayEquals(expected, actual); +} + +TEST_F(MapSharedShreddingFileReaderTest, TestSelectedKeysStructProjectionFromDefaultMap) { + auto map_type = arrow::internal::checked_pointer_cast( + arrow::map(arrow::utf8(), arrow::field("value", arrow::int64()))); + auto file_schema = + arrow::schema({arrow::field("id", arrow::int32()), arrow::field("tags", map_type)}); + auto file_array = + arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(file_schema->fields()), + R"([ + [1, [["a", 10], ["c", null]]], + [2, [["b", 20]]], + [3, null] + ])") + .ValueOrDie(); + auto mock_reader = std::make_unique( + file_array, arrow::struct_(file_schema->fields()), /*read_batch_size=*/10); + mock_reader->EnableRandomizeBatchSize(false); + + auto selected_type = arrow::struct_( + {arrow::field("a", arrow::int64()), arrow::field("missing", arrow::int64())}); + auto selected_field = + arrow::field("tags", selected_type, /*nullable=*/true, + arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,missing"})); + ASSERT_OK_AND_ASSIGN(auto field_read_plan, + MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan( + file_schema->field(1), selected_field)); + std::map> contexts; + contexts.emplace("tags", std::move(field_read_plan)); + auto reader = std::make_unique(std::move(mock_reader), + std::move(contexts), pool_); + + auto read_schema = + ExportSchema(arrow::schema({arrow::field("id", arrow::int32()), selected_field})); + ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr, + /*selection_bitmap=*/std::nullopt)); + ASSERT_OK_AND_ASSIGN(auto actual, ReadResultCollector::CollectResult(reader.get())); + + auto expected_type = arrow::struct_({arrow::field("id", arrow::int32()), selected_field}); + std::shared_ptr expected; + ASSERT_TRUE(arrow::ipc::internal::json::ChunkedArrayFromJSON(expected_type, {R"([ + [1, [10, null]], + [2, [null, null]], + [3, null] + ])"}, + &expected) + .ok()); + AssertChunkedArrayEquals(expected, actual); +} + +TEST_F(MapSharedShreddingFileReaderTest, TestInvalidSelectedKeysStructProjection) { + auto file_map_field = arrow::field("tags", arrow::map(arrow::utf8(), arrow::int64())); + auto mismatched_count_field = + arrow::field("tags", arrow::struct_({arrow::field("a", arrow::int64())}), /*nullable=*/true, + arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,b"})); + ASSERT_NOK_WITH_MSG(MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan( + mismatched_count_field, TagsMeta()), + "metadata size 2 does not match STRUCT field count 1"); + ASSERT_NOK_WITH_MSG(MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan( + file_map_field, mismatched_count_field), + "metadata size 2 does not match STRUCT field count 1"); + + auto mismatched_type_field = arrow::field( + "tags", + arrow::struct_({arrow::field("a", arrow::int64()), arrow::field("b", arrow::utf8())}), + /*nullable=*/true, arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,b"})); + ASSERT_NOK_WITH_MSG(MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan( + mismatched_type_field, TagsMeta()), + "must have the same value type"); + ASSERT_NOK_WITH_MSG(MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan( + file_map_field, mismatched_type_field), + "must have the same value type"); +} + TEST_F(MapSharedShreddingFileReaderTest, TestPartialExistSelectedKeys) { ASSERT_OK_AND_ASSIGN(auto reader, CreateReader(/*physical_array=*/nullptr, /*physical_schema=*/nullptr, diff --git a/src/paimon/common/data/shredding/map_shared_shredding_schema_utils.cpp b/src/paimon/common/data/shredding/map_shared_shredding_schema_utils.cpp index c8dc54128..36978d04f 100644 --- a/src/paimon/common/data/shredding/map_shared_shredding_schema_utils.cpp +++ b/src/paimon/common/data/shredding/map_shared_shredding_schema_utils.cpp @@ -19,15 +19,96 @@ #include "paimon/data/shredding/map_shared_shredding_schema_utils.h" +#include +#include +#include + #include "arrow/c/bridge.h" #include "arrow/type.h" #include "arrow/util/key_value_metadata.h" #include "fmt/format.h" #include "paimon/common/data/shredding/map_shared_shredding_utils.h" +#include "paimon/common/types/data_field.h" #include "paimon/common/utils/arrow/status_utils.h" namespace paimon { +class MapSharedShreddingAccessBuilder::Impl { + public: + Impl(const std::shared_ptr& _map_field, + const std::shared_ptr& _map_type) + : map_field(_map_field), map_type(_map_type) {} + + std::shared_ptr map_field; + std::shared_ptr map_type; + std::vector keys; + std::unordered_set unique_keys; +}; + +MapSharedShreddingAccessBuilder::~MapSharedShreddingAccessBuilder() = default; + +MapSharedShreddingAccessBuilder::MapSharedShreddingAccessBuilder(std::unique_ptr&& impl) + : impl_(std::move(impl)) {} + +Result> MapSharedShreddingAccessBuilder::Create( + struct ArrowSchema* map_field) { + if (!map_field) { + return Status::Invalid("MAP field is null"); + } + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr field, + arrow::ImportField(map_field)); + if (field->type()->id() != arrow::Type::MAP) { + return Status::Invalid( + fmt::format("MapSharedShreddingAccessBuilder requires MAP field, got {}", + field->type()->ToString())); + } + auto map_type = arrow::internal::checked_pointer_cast(field->type()); + if (map_type->key_type()->id() != arrow::Type::STRING) { + return Status::Invalid(fmt::format( + "MapSharedShreddingAccessBuilder only supports MAP with STRING keys, got {}", + map_type->key_type()->ToString())); + } + auto impl = std::make_unique(field, map_type); + return std::unique_ptr( + new MapSharedShreddingAccessBuilder(std::move(impl))); +} + +Status MapSharedShreddingAccessBuilder::AddKey(const std::string& key) { + if (key.find(',') != std::string::npos) { + return Status::Invalid( + fmt::format("selected MAP key {} must not contain the ',' delimiter", key)); + } + if (!impl_->unique_keys.insert(key).second) { + return Status::Invalid(fmt::format("selected MAP key must not be duplicated: {}", key)); + } + impl_->keys.push_back(key); + return Status::OK(); +} + +Result> MapSharedShreddingAccessBuilder::Build() const { + if (impl_->keys.empty()) { + return Status::Invalid( + "shared shredding MAP selected-key projection needs at least one key"); + } + arrow::FieldVector fields; + fields.reserve(impl_->keys.size()); + std::string encoded_keys; + for (size_t i = 0; i < impl_->keys.size(); ++i) { + if (i != 0) { + encoded_keys.push_back(','); + } + encoded_keys.append(impl_->keys[i]); + fields.push_back(arrow::field(impl_->keys[i], impl_->map_type->item_type(), + /*nullable=*/true)); + } + auto metadata = arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {encoded_keys}); + auto access_field = impl_->map_field->WithType(arrow::struct_(std::move(fields))) + ->WithMetadata(std::move(metadata)); + auto field = std::make_unique(); + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportField(*access_field, field.get())); + return field; +} + Result> MapSharedShreddingSchemaUtils::LogicalToPhysicalSchema( std::unique_ptr<::ArrowSchema> logical_schema, const std::map& field_to_num_columns) { diff --git a/src/paimon/common/data/shredding/map_shared_shredding_schema_utils_test.cpp b/src/paimon/common/data/shredding/map_shared_shredding_schema_utils_test.cpp index e421bc3f6..bfe26a372 100644 --- a/src/paimon/common/data/shredding/map_shared_shredding_schema_utils_test.cpp +++ b/src/paimon/common/data/shredding/map_shared_shredding_schema_utils_test.cpp @@ -27,9 +27,85 @@ #include "arrow/util/key_value_metadata.h" #include "gtest/gtest.h" #include "paimon/common/data/shredding/map_shared_shredding_utils.h" +#include "paimon/common/types/data_field.h" #include "paimon/testing/utils/testharness.h" namespace paimon::test { +namespace { + +std::unique_ptr ExportField(const std::shared_ptr& field) { + auto c_field = std::make_unique(); + EXPECT_TRUE(arrow::ExportField(*field, c_field.get()).ok()); + return c_field; +} + +} // namespace + +TEST(MapSharedShreddingAccessBuilderTest, BuildSelectedKeysField) { + auto original_metadata = + arrow::KeyValueMetadata::Make({DataField::FIELD_ID, DataField::DESCRIPTION, "custom.key"}, + {"7", "original description", "custom.value"}); + auto map_type = arrow::map(arrow::utf8(), arrow::field("value", arrow::int64(), false)); + auto map_field = arrow::field("attributes", map_type, /*nullable=*/false, original_metadata); + ASSERT_OK_AND_ASSIGN(std::unique_ptr builder, + MapSharedShreddingAccessBuilder::Create(ExportField(map_field).get())); + ASSERT_OK(builder->AddKey("age")); + ASSERT_OK(builder->AddKey("score")); + + ASSERT_OK_AND_ASSIGN(std::unique_ptr c_field, builder->Build()); + auto imported_field = arrow::ImportField(c_field.get()); + ASSERT_TRUE(imported_field.ok()); + std::shared_ptr field = imported_field.ValueOrDie(); + ASSERT_EQ(field->name(), "attributes"); + ASSERT_EQ(field->type()->id(), arrow::Type::STRUCT); + ASSERT_FALSE(field->nullable()); + + auto struct_type = arrow::internal::checked_pointer_cast(field->type()); + ASSERT_EQ(struct_type->num_fields(), 2); + ASSERT_EQ(struct_type->field(0)->name(), "age"); + ASSERT_EQ(struct_type->field(1)->name(), "score"); + ASSERT_TRUE(struct_type->field(0)->type()->Equals(arrow::int64())); + ASSERT_TRUE(struct_type->field(1)->type()->Equals(arrow::int64())); + ASSERT_TRUE(struct_type->field(0)->nullable()); + ASSERT_TRUE(struct_type->field(1)->nullable()); + ASSERT_FALSE(field->metadata()->Contains(DataField::FIELD_ID)); + ASSERT_FALSE(field->metadata()->Contains(DataField::DESCRIPTION)); + ASSERT_FALSE(field->metadata()->Contains("custom.key")); + ASSERT_TRUE(field->metadata()->Contains(DataField::MAP_SELECTED_KEYS)); + ASSERT_EQ(field->metadata()->Get(DataField::MAP_SELECTED_KEYS).ValueOrDie(), "age,score"); +} + +TEST(MapSharedShreddingAccessBuilderTest, RejectInvalidKeys) { + { + auto map_field = arrow::field("attributes", arrow::map(arrow::utf8(), arrow::int64())); + ASSERT_OK_AND_ASSIGN(std::unique_ptr builder, + MapSharedShreddingAccessBuilder::Create(ExportField(map_field).get())); + ASSERT_NOK_WITH_MSG(builder->Build(), "at least one key"); + } + { + auto map_field = arrow::field("attributes", arrow::map(arrow::utf8(), arrow::int64())); + ASSERT_OK_AND_ASSIGN(std::unique_ptr builder, + MapSharedShreddingAccessBuilder::Create(ExportField(map_field).get())); + ASSERT_OK(builder->AddKey("a")); + ASSERT_NOK_WITH_MSG(builder->AddKey("a"), "must not be duplicated"); + } + { + auto map_field = arrow::field("attributes", arrow::map(arrow::utf8(), arrow::int64())); + ASSERT_OK_AND_ASSIGN(std::unique_ptr builder, + MapSharedShreddingAccessBuilder::Create(ExportField(map_field).get())); + ASSERT_NOK_WITH_MSG(builder->AddKey("a,b"), "must not contain the ',' delimiter"); + } +} + +TEST(MapSharedShreddingAccessBuilderTest, RejectInvalidMapField) { + ASSERT_NOK_WITH_MSG(MapSharedShreddingAccessBuilder::Create(nullptr), "MAP field is null"); + ASSERT_NOK_WITH_MSG(MapSharedShreddingAccessBuilder::Create( + ExportField(arrow::field("v", arrow::int64())).get()), + "requires MAP field"); + auto non_string_map = arrow::field("attributes", arrow::map(arrow::int32(), arrow::int64())); + ASSERT_NOK_WITH_MSG(MapSharedShreddingAccessBuilder::Create(ExportField(non_string_map).get()), + "only supports MAP with STRING keys"); +} TEST(MapSharedShreddingSchemaUtilsTest, AttachMetadataToSchemaBasic) { MapSharedShreddingFieldMeta tags_meta; diff --git a/src/paimon/core/io/field_mapping_reader.cpp b/src/paimon/core/io/field_mapping_reader.cpp index 3bd04b49b..947853c3c 100644 --- a/src/paimon/core/io/field_mapping_reader.cpp +++ b/src/paimon/core/io/field_mapping_reader.cpp @@ -51,6 +51,18 @@ Result FieldMappingReader::HasMapSelectedKeysRecursively( return false; } auto type_id = read_field->type()->id(); + if (NestedProjectionUtils::IsMapSharedShreddingAccessField(read_field)) { + PAIMON_ASSIGN_OR_RAISE(std::vector selected_keys, + NestedProjectionUtils::GetMapSelectedKeys(read_field)); + auto read_struct = + arrow::internal::checked_pointer_cast(read_field->type()); + if (selected_keys.size() != static_cast(read_struct->num_fields())) { + return Status::Invalid(fmt::format( + "selected-key metadata size {} does not match STRUCT field count {} for {}", + selected_keys.size(), read_struct->num_fields(), read_field->name())); + } + return true; + } if (type_id == arrow::Type::MAP) { PAIMON_ASSIGN_OR_RAISE(std::vector selected_keys, NestedProjectionUtils::GetMapSelectedKeys(read_field)); @@ -75,6 +87,11 @@ Result> FieldMappingReader::FilterMapSelectedKeysR } auto type_id = read_field->type()->id(); + if (NestedProjectionUtils::IsMapSharedShreddingAccessField(read_field)) { + // The shared-shredding wrapper (including its default MAP fallback) has already + // materialized this projection as a STRUCT. + return array; + } if (type_id == arrow::Type::MAP) { PAIMON_ASSIGN_OR_RAISE(std::vector selected_keys, NestedProjectionUtils::GetMapSelectedKeys(read_field)); diff --git a/src/paimon/core/operation/abstract_split_read.cpp b/src/paimon/core/operation/abstract_split_read.cpp index 0850d943f..2c904ca09 100644 --- a/src/paimon/core/operation/abstract_split_read.cpp +++ b/src/paimon/core/operation/abstract_split_read.cpp @@ -256,8 +256,7 @@ AbstractSplitRead::ApplySharedShreddingReaderIfNeeded( file_reader->GetFileSchema()); PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr file_arrow_schema, arrow::ImportSchema(file_schema.get())); - std::map - shared_shredding_name_to_context; + std::map> field_read_plans; for (const auto& read_field : read_schema->fields()) { const auto& field_name = read_field->name(); auto file_field = file_arrow_schema->GetFieldByName(field_name); @@ -267,35 +266,41 @@ AbstractSplitRead::ApplySharedShreddingReaderIfNeeded( } std::shared_ptr metadata = std::const_pointer_cast(file_field->metadata()); - if (!MapSharedShreddingUtils::HasShreddingMetadata(metadata)) { - // not a map shared shredding field + bool is_shared_shredding_file = MapSharedShreddingUtils::HasShreddingMetadata(metadata); + bool is_shared_shredding_map_access = + NestedProjectionUtils::IsMapSharedShreddingAccessField(read_field); + if (!is_shared_shredding_file && !is_shared_shredding_map_access) { + // Neither a shared-shredding file field nor a selected-key STRUCT projection. continue; } - // get meta - PAIMON_ASSIGN_OR_RAISE(MapSharedShreddingFieldMeta meta, - MapSharedShreddingUtils::DeserializeMetadata(metadata)); - // get selected_keys - PAIMON_ASSIGN_OR_RAISE(std::vector selected_keys, - NestedProjectionUtils::GetMapSelectedKeys(read_field)); - if (selected_keys.empty()) { - // select all keys - selected_keys.reserve(meta.name_to_id.size()); - for (const auto& [key_name, _] : meta.name_to_id) { - selected_keys.push_back(key_name); + + std::unique_ptr field_read_plan; + if (is_shared_shredding_map_access) { + if (is_shared_shredding_file) { + PAIMON_ASSIGN_OR_RAISE(MapSharedShreddingFieldMeta meta, + MapSharedShreddingUtils::DeserializeMetadata(metadata)); + PAIMON_ASSIGN_OR_RAISE( + field_read_plan, + MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(read_field, meta)); + } else { + PAIMON_ASSIGN_OR_RAISE(field_read_plan, + MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan( + file_field, read_field)); } + } else { + PAIMON_ASSIGN_OR_RAISE(MapSharedShreddingFieldMeta meta, + MapSharedShreddingUtils::DeserializeMetadata(metadata)); + PAIMON_ASSIGN_OR_RAISE(field_read_plan, + MapFieldReadPlanFactory::CreateMapReadPlan(read_field, meta)); } - // get map type - auto map_type = arrow::internal::checked_pointer_cast(read_field->type()); - shared_shredding_name_to_context.emplace( - field_name, - MapSharedShreddingFileReader::SharedShreddingContext(meta, selected_keys, map_type)); + field_read_plans.emplace(field_name, std::move(field_read_plan)); PAIMON_ASSIGN_OR_RAISE(int32_t field_id, NestedProjectionUtils::GetPaimonFieldId(read_field)); handled_shared_shredding_field_ids.insert(field_id); } - if (!shared_shredding_name_to_context.empty()) { + if (!field_read_plans.empty()) { file_reader = std::make_unique( - std::move(file_reader), std::move(shared_shredding_name_to_context), pool_); + std::move(file_reader), std::move(field_read_plans), pool_); } return std::make_pair(std::move(file_reader), std::move(handled_shared_shredding_field_ids)); } diff --git a/src/paimon/core/operation/internal_read_context.cpp b/src/paimon/core/operation/internal_read_context.cpp index e8886b426..7c4e2ecdd 100644 --- a/src/paimon/core/operation/internal_read_context.cpp +++ b/src/paimon/core/operation/internal_read_context.cpp @@ -31,6 +31,7 @@ #include "paimon/common/table/special_fields.h" #include "paimon/common/types/data_field.h" #include "paimon/common/utils/arrow/status_utils.h" +#include "paimon/core/options/map_storage_layout.h" #include "paimon/core/schema/arrow_schema_validator.h" #include "paimon/core/utils/nested_projection_utils.h" #include "paimon/status.h" @@ -50,6 +51,31 @@ Result> InternalReadContext::AlignReadFieldWithTab return table_field->WithType(read_field->type()); } + if (table_field->type()->id() == arrow::Type::MAP && + NestedProjectionUtils::IsMapSharedShreddingAccessField(read_field)) { + auto table_map = arrow::internal::checked_pointer_cast(table_field->type()); + if (table_map->key_type()->id() != arrow::Type::STRING) { + return Status::Invalid(fmt::format( + "Selected-key MAP pushdown only supports string MAP keys for field '{}'", + table_field->name())); + } + PAIMON_RETURN_NOT_OK( + NestedProjectionUtils::ValidateMapSharedShreddingAccessField(read_field).status()); + auto read_struct = + arrow::internal::checked_pointer_cast(read_field->type()); + const auto& selected_value_type = read_struct->field(0)->type(); + if (!selected_value_type->Equals(table_map->item_type())) { + return Status::Invalid(fmt::format( + "Selected-key MAP pushdown does not support pruning MAP value fields for " + "'{}': selected type {} vs MAP value type {}", + table_field->name(), selected_value_type->ToString(), + table_map->item_type()->ToString())); + } + auto aligned_field = table_field->WithType(read_field->type()); + return DataField::MergeFieldMetadataByWhitelist(aligned_field, read_field, + kReadMetadataWhitelist); + } + if (read_field->type()->id() != table_field->type()->id()) { return Status::Invalid(fmt::format( "Read schema field '{}' type {} does not match table field type {}", read_field->name(), @@ -198,6 +224,16 @@ Result> InternalReadContext::Create( } PAIMON_ASSIGN_OR_RAISE(DataField table_field, table_schema->GetField(read_field->name())); + if (NestedProjectionUtils::IsMapSharedShreddingAccessField(read_field)) { + PAIMON_ASSIGN_OR_RAISE(MapStorageLayout layout, + core_options.GetMapStorageLayout(table_field.Name())); + if (layout != MapStorageLayout::SHARED_SHREDDING) { + return Status::Invalid(fmt::format( + "Selected-key MAP pushdown only supports top-level shared-shredding MAP " + "field: {}", + table_field.Name())); + } + } PAIMON_ASSIGN_OR_RAISE( std::shared_ptr aligned_field, AlignReadFieldWithTableFieldIds(read_field, table_field.ArrowField())); diff --git a/src/paimon/core/operation/internal_read_context_test.cpp b/src/paimon/core/operation/internal_read_context_test.cpp index 30ba77b70..286188570 100644 --- a/src/paimon/core/operation/internal_read_context_test.cpp +++ b/src/paimon/core/operation/internal_read_context_test.cpp @@ -25,6 +25,7 @@ #include "paimon/common/table/special_fields.h" #include "paimon/common/types/data_field.h" #include "paimon/core/schema/schema_manager.h" +#include "paimon/data/shredding/map_shared_shredding_schema_utils.h" #include "paimon/defs.h" #include "paimon/fs/local/local_file_system.h" #include "paimon/status.h" @@ -303,4 +304,34 @@ TEST(InternalReadContext, TestProjectedSchemaMetadataWhitelist) { ASSERT_FALSE(custom_metadata_result.ok()); } +TEST(InternalReadContext, TestMapSharedShreddingAccessRequiresSharedShreddingLayout) { + auto map_field = arrow::field("tags", arrow::map(arrow::utf8(), arrow::int64())); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr unique_table_schema, + TableSchema::Create(/*schema_id=*/0, arrow::schema({map_field}), + /*partition_keys=*/{}, /*primary_keys=*/{}, /*options=*/{})); + std::shared_ptr table_schema = std::move(unique_table_schema); + + auto c_map_field = std::make_unique(); + ASSERT_TRUE(arrow::ExportField(*map_field, c_map_field.get()).ok()); + ASSERT_OK_AND_ASSIGN(std::unique_ptr access_builder, + MapSharedShreddingAccessBuilder::Create(c_map_field.get())); + ASSERT_OK(access_builder->AddKey("a")); + ASSERT_OK_AND_ASSIGN(std::unique_ptr c_access_field, access_builder->Build()); + auto imported_access_field = arrow::ImportField(c_access_field.get()); + ASSERT_TRUE(imported_access_field.ok()); + std::shared_ptr access_field = imported_access_field.ValueOrDie(); + + auto c_read_schema = std::make_unique(); + ASSERT_TRUE(arrow::ExportSchema(*arrow::schema({access_field}), c_read_schema.get()).ok()); + ReadContextBuilder context_builder("/tmp/unused-table-path"); + context_builder.SetReadSchema(std::move(c_read_schema)); + ASSERT_OK_AND_ASSIGN(auto unique_read_context, context_builder.Finish()); + std::shared_ptr read_context = std::move(unique_read_context); + + ASSERT_NOK_WITH_MSG( + InternalReadContext::Create(read_context, table_schema, table_schema->Options()), + "Selected-key MAP pushdown only supports top-level shared-shredding MAP field: tags"); +} + } // namespace paimon::test diff --git a/src/paimon/core/utils/field_mapping.cpp b/src/paimon/core/utils/field_mapping.cpp index df3791e33..c9b70d476 100644 --- a/src/paimon/core/utils/field_mapping.cpp +++ b/src/paimon/core/utils/field_mapping.cpp @@ -108,9 +108,16 @@ Result FieldMappingBuilder::CreateExistFieldInfo( // Recursively prune nested types in data_field to match read_field's // projection. For atomic types this is a no-op. - PAIMON_ASSIGN_OR_RAISE( - std::optional> pruned_type, - NestedProjectionUtils::PruneDataType(read_field.Type(), data_field.Type())); + std::optional> pruned_type; + if (NestedProjectionUtils::IsMapSharedShreddingAccessField(read_field.ArrowField())) { + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr map_access_data_type, + NestedProjectionUtils::BuildMapSharedShreddingAccessDataType( + read_field.ArrowField(), data_field.Type())); + pruned_type = std::move(map_access_data_type); + } else { + PAIMON_ASSIGN_OR_RAISE(pruned_type, NestedProjectionUtils::PruneDataType( + read_field.Type(), data_field.Type())); + } if (!pruned_type.has_value()) { // All sub-fields pruned away — treat as non-existent. continue; diff --git a/src/paimon/core/utils/nested_projection_utils.cpp b/src/paimon/core/utils/nested_projection_utils.cpp index 97c08e9a0..3826ae81d 100644 --- a/src/paimon/core/utils/nested_projection_utils.cpp +++ b/src/paimon/core/utils/nested_projection_utils.cpp @@ -446,79 +446,104 @@ Result> NestedProjectionUtils::GetMapSelectedKeys( return result; } -namespace { +bool NestedProjectionUtils::IsMapSharedShreddingAccessField( + const std::shared_ptr& field) { + if (field->type()->id() != arrow::Type::STRUCT || !field->HasMetadata() || !field->metadata()) { + return false; + } + return field->metadata()->Contains(DataField::MAP_SELECTED_KEYS); +} -struct MapKeyAccessor { - std::shared_ptr string_keys; - std::shared_ptr dict_keys; - std::shared_ptr dict_values; - std::shared_ptr dict_large_values; -}; +Result> NestedProjectionUtils::ValidateMapSharedShreddingAccessField( + const std::shared_ptr& field) { + if (field->type()->id() != arrow::Type::STRUCT) { + return Status::Invalid( + fmt::format("selected-key MAP field {} is not a STRUCT", field->name())); + } + auto struct_type = arrow::internal::checked_pointer_cast(field->type()); + PAIMON_ASSIGN_OR_RAISE(std::vector selected_keys, GetMapSelectedKeys(field)); + if (struct_type->num_fields() == 0 || + selected_keys.size() != static_cast(struct_type->num_fields())) { + return Status::Invalid( + fmt::format("selected-key metadata size {} does not match STRUCT field count {} for {}", + selected_keys.size(), struct_type->num_fields(), field->name())); + } + const auto& value_type = struct_type->field(0)->type(); + for (int32_t i = 1; i < struct_type->num_fields(); ++i) { + if (!struct_type->field(i)->type()->Equals(value_type)) { + return Status::Invalid(fmt::format( + "selected-key MAP fields must have the same value type, but {} and {} differ", + value_type->ToString(), struct_type->field(i)->type()->ToString())); + } + } + return selected_keys; +} + +Result> +NestedProjectionUtils::BuildMapSharedShreddingAccessDataType( + const std::shared_ptr& read_field, + const std::shared_ptr& data_type) { + if (!IsMapSharedShreddingAccessField(read_field)) { + return Status::Invalid( + fmt::format("field {} is not a selected-key MAP projection", read_field->name())); + } + if (data_type->id() != arrow::Type::MAP) { + return Status::Invalid( + fmt::format("selected-key MAP projection {} requires MAP data type, got {}", + read_field->name(), data_type->ToString())); + } + PAIMON_RETURN_NOT_OK(ValidateMapSharedShreddingAccessField(read_field).status()); + auto read_struct = arrow::internal::checked_pointer_cast(read_field->type()); + auto data_map = arrow::internal::checked_pointer_cast(data_type); + arrow::FieldVector data_children; + data_children.reserve(read_struct->num_fields()); + for (const auto& read_child : read_struct->fields()) { + data_children.push_back(read_child->WithType(data_map->item_type())); + } + return arrow::struct_(std::move(data_children)); +} -Result BuildMapKeyAccessor(const std::shared_ptr& key_array) { - MapKeyAccessor accessor; +Result NestedProjectionUtils::GetMapKeyViewAt( + const std::shared_ptr& key_array, int64_t entry_idx) { + if (key_array->IsNull(entry_idx)) { + return Status::Invalid("selected-key MAP read found null MAP key at entry " + + std::to_string(entry_idx)); + } if (key_array->type_id() == arrow::Type::STRING) { - accessor.string_keys = std::static_pointer_cast(key_array); - return accessor; + return arrow::internal::checked_pointer_cast(key_array)->GetView( + entry_idx); } if (key_array->type_id() == arrow::Type::DICTIONARY) { - auto dict_type = std::static_pointer_cast(key_array->type()); + auto dict_type = + arrow::internal::checked_pointer_cast(key_array->type()); if (dict_type->value_type()->id() != arrow::Type::STRING && dict_type->value_type()->id() != arrow::Type::LARGE_STRING) { return Status::Invalid( - fmt::format("FilterMapArrayBySelectedKeys only supports string keys or " + fmt::format("selected-key MAP read only supports string keys or " "dictionary keys, got {}", key_array->type()->ToString())); } - accessor.dict_keys = std::static_pointer_cast(key_array); + auto dict_keys = arrow::internal::checked_pointer_cast(key_array); + int64_t dict_idx = dict_keys->GetValueIndex(entry_idx); + const auto& dictionary = dict_keys->dictionary(); + if (dictionary->IsNull(dict_idx)) { + return Status::Invalid( + "selected-key MAP read found null dictionary MAP key at dictionary index " + + std::to_string(dict_idx)); + } if (dict_type->value_type()->id() == arrow::Type::STRING) { - accessor.dict_values = - std::static_pointer_cast(accessor.dict_keys->dictionary()); - } else { - accessor.dict_large_values = - std::static_pointer_cast(accessor.dict_keys->dictionary()); + return arrow::internal::checked_pointer_cast(dictionary) + ->GetView(dict_idx); } - return accessor; + return arrow::internal::checked_pointer_cast(dictionary) + ->GetView(dict_idx); } return Status::Invalid( - fmt::format("FilterMapArrayBySelectedKeys only supports string keys or " + fmt::format("selected-key MAP read only supports string keys or " "dictionary keys, got {}", key_array->type()->ToString())); } -Result GetMapKeyViewAt(const MapKeyAccessor& accessor, int64_t entry_idx) { - if (accessor.string_keys) { - if (accessor.string_keys->IsNull(entry_idx)) { - return Status::Invalid("FilterMapArrayBySelectedKeys found null map key at entry " + - std::to_string(entry_idx)); - } - return accessor.string_keys->GetView(entry_idx); - } - - if (accessor.dict_keys->IsNull(entry_idx)) { - return Status::Invalid("FilterMapArrayBySelectedKeys found null map key at entry " + - std::to_string(entry_idx)); - } - int64_t dict_idx = accessor.dict_keys->GetValueIndex(entry_idx); - if (accessor.dict_values) { - if (accessor.dict_values->IsNull(dict_idx)) { - return Status::Invalid( - "FilterMapArrayBySelectedKeys found null dictionary map key at dictionary index " + - std::to_string(dict_idx)); - } - return accessor.dict_values->GetView(dict_idx); - } - - if (accessor.dict_large_values->IsNull(dict_idx)) { - return Status::Invalid( - "FilterMapArrayBySelectedKeys found null dictionary map key at dictionary index " + - std::to_string(dict_idx)); - } - return accessor.dict_large_values->GetView(dict_idx); -} - -} // namespace - Result> NestedProjectionUtils::FilterMapArrayBySelectedKeys( const std::shared_ptr& array, const std::vector& selected_keys, arrow::MemoryPool* pool) { @@ -534,12 +559,11 @@ Result> NestedProjectionUtils::FilterMapArrayBySel "FilterMapArrayBySelectedKeys requires map array, got {}", array->type()->ToString())); } - auto map_array = std::static_pointer_cast(array); - auto map_type = std::static_pointer_cast(array->type()); + auto map_array = arrow::internal::checked_pointer_cast(array); + auto map_type = arrow::internal::checked_pointer_cast(array->type()); assert(map_array && map_type); auto key_array = map_array->keys(); - PAIMON_ASSIGN_OR_RAISE(MapKeyAccessor key_accessor, BuildMapKeyAccessor(key_array)); auto values_array = map_array->items(); int64_t num_maps = map_array->length(); @@ -575,7 +599,7 @@ Result> NestedProjectionUtils::FilterMapArrayBySel for (const auto& selected_key : selected_keys) { for (int64_t entry_idx = start; entry_idx < end; ++entry_idx) { PAIMON_ASSIGN_OR_RAISE(std::string_view key_view, - GetMapKeyViewAt(key_accessor, entry_idx)); + GetMapKeyViewAt(key_array, entry_idx)); if (key_view == selected_key) { PAIMON_RETURN_NOT_OK_FROM_ARROW(key_builder->Append( key_view.data(), static_cast(key_view.size()))); diff --git a/src/paimon/core/utils/nested_projection_utils.h b/src/paimon/core/utils/nested_projection_utils.h index b0ab8fc2f..59e7d95f2 100644 --- a/src/paimon/core/utils/nested_projection_utils.h +++ b/src/paimon/core/utils/nested_projection_utils.h @@ -79,6 +79,27 @@ class PAIMON_EXPORT NestedProjectionUtils { static Result> GetMapSelectedKeys( const std::shared_ptr& field); + /// @return true when `field` is a selected-key MAP projection: a STRUCT carrying + /// `paimon.map.selected-keys` metadata. + static bool IsMapSharedShreddingAccessField(const std::shared_ptr& field); + + /// Validates a selected-key MAP projection and returns its selected keys. The field must be a + /// non-empty STRUCT, its metadata key count must match its child count, and all children must + /// have the same value type. + static Result> ValidateMapSharedShreddingAccessField( + const std::shared_ptr& field); + + /// Rewrites a selected-key STRUCT projection to use the data file's complete MAP value type + /// for every child before materialization. Cpp paimon does not support schema evolution for + /// for field inside the MAP value. + static Result> BuildMapSharedShreddingAccessDataType( + const std::shared_ptr& read_field, + const std::shared_ptr& data_type); + + /// Returns a string view for a MAP key stored as string or dictionary. + static Result GetMapKeyViewAt(const std::shared_ptr& key_array, + int64_t entry_idx); + /// Filter a MapArray so that only entries whose key is in `selected_keys` are kept. /// Supports string keys and dictionary keys. /// The output map entry order follows diff --git a/src/paimon/core/utils/nested_projection_utils_test.cpp b/src/paimon/core/utils/nested_projection_utils_test.cpp index 7905b2ce9..570a5321d 100644 --- a/src/paimon/core/utils/nested_projection_utils_test.cpp +++ b/src/paimon/core/utils/nested_projection_utils_test.cpp @@ -547,6 +547,62 @@ TEST(NestedProjectionUtilsTest, GetMapSelectedKeysDuplicateKey) { "Duplicate selected key 'a'"); } +// ============== MapSharedShreddingAccessField ============== + +TEST(NestedProjectionUtilsTest, IsMapSharedShreddingAccessField) { + auto metadata = arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,b"}); + auto access_type = + arrow::struct_({arrow::field("a", arrow::int64()), arrow::field("b", arrow::int64())}); + + ASSERT_TRUE(NestedProjectionUtils::IsMapSharedShreddingAccessField( + arrow::field("tags", access_type, /*nullable=*/true, metadata))); + ASSERT_FALSE( + NestedProjectionUtils::IsMapSharedShreddingAccessField(arrow::field("tags", access_type))); + ASSERT_FALSE(NestedProjectionUtils::IsMapSharedShreddingAccessField(arrow::field( + "tags", arrow::map(arrow::utf8(), arrow::int64()), /*nullable=*/true, metadata))); +} + +TEST(NestedProjectionUtilsTest, BuildMapSharedShreddingAccessDataType) { + auto read_type = arrow::struct_({ + arrow::field("a", arrow::int64(), /*nullable=*/true), + arrow::field("b", arrow::int64(), /*nullable=*/true), + }); + auto read_metadata = arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,b"}); + auto read_field = arrow::field("tags", read_type, /*nullable=*/true, std::move(read_metadata)); + auto data_value_type = arrow::int32(); + auto data_type = arrow::map(arrow::utf8(), data_value_type); + + ASSERT_OK_AND_ASSIGN( + std::shared_ptr result, + NestedProjectionUtils::BuildMapSharedShreddingAccessDataType(read_field, data_type)); + auto result_struct = arrow::internal::checked_pointer_cast(result); + ASSERT_EQ(result_struct->num_fields(), 2); + ASSERT_EQ(result_struct->field(0)->name(), "a"); + ASSERT_EQ(result_struct->field(1)->name(), "b"); + ASSERT_TRUE(result_struct->field(0)->type()->Equals(data_value_type)); + ASSERT_TRUE(result_struct->field(1)->type()->Equals(data_value_type)); + ASSERT_TRUE(result_struct->field(0)->nullable()); + ASSERT_TRUE(result_struct->field(1)->nullable()); +} + +TEST(NestedProjectionUtilsTest, BuildMapSharedShreddingAccessDataTypeInvalidInput) { + auto access_metadata = arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,b"}); + auto access_field = arrow::field("tags", arrow::struct_({arrow::field("a", arrow::int64())}), + /*nullable=*/true, access_metadata); + + ASSERT_NOK_WITH_MSG( + NestedProjectionUtils::BuildMapSharedShreddingAccessDataType( + arrow::field("tags", arrow::struct_({arrow::field("a", arrow::int64())})), + arrow::map(arrow::utf8(), arrow::int64())), + "is not a selected-key MAP projection"); + ASSERT_NOK_WITH_MSG( + NestedProjectionUtils::BuildMapSharedShreddingAccessDataType(access_field, arrow::int64()), + "requires MAP data type"); + ASSERT_NOK_WITH_MSG(NestedProjectionUtils::BuildMapSharedShreddingAccessDataType( + access_field, arrow::map(arrow::utf8(), arrow::int64())), + "metadata size 2 does not match STRUCT field count 1"); +} + // ============== FilterMapArrayBySelectedKeys ============== class NestedProjectionUtilsMapArrayTest : public ::testing::Test { diff --git a/test/inte/write_and_read_inte_test.cpp b/test/inte/write_and_read_inte_test.cpp index 9b8f94446..bdb0e2db8 100644 --- a/test/inte/write_and_read_inte_test.cpp +++ b/test/inte/write_and_read_inte_test.cpp @@ -40,6 +40,7 @@ #include "paimon/core/io/data_file_meta.h" #include "paimon/core/schema/schema_manager.h" #include "paimon/core/table/source/data_split_impl.h" +#include "paimon/data/shredding/map_shared_shredding_schema_utils.h" #include "paimon/defs.h" #include "paimon/file_store_commit.h" #include "paimon/file_store_write.h" @@ -167,6 +168,23 @@ class WriteAndReadInteTest return std::make_shared(expected)->Equals(actual); } + Result> BuildMapSharedShreddingAccessField( + const std::shared_ptr& map_field, + const std::vector& selected_keys) const { + auto c_map_field = std::make_unique(); + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportField(*map_field, c_map_field.get())); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr access_builder, + MapSharedShreddingAccessBuilder::Create(c_map_field.get())); + for (const auto& selected_key : selected_keys) { + PAIMON_RETURN_NOT_OK(access_builder->AddKey(selected_key)); + } + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr c_access_field, + access_builder->Build()); + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr access_field, + arrow::ImportField(c_access_field.get())); + return access_field; + } + Result> InnerScan( const std::map& options) const { std::string table_path = PathUtil::JoinPath(test_dir_, "foo.db/bar"); @@ -2330,6 +2348,21 @@ TEST_P(WriteAndReadInteTest, TestMapSharedShreddingReadAfterRenameColumn) { [0, 2, [["c", 21]]] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields_v1[1], {"b", "a"})); + auto read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options_v1, read_schema, expected_type, + R"([ + [0, 1, [12, 11]], + [0, 2, [null, null]] + ])")); + ASSERT_TRUE(success); } TEST_P(WriteAndReadInteTest, TestSharedShreddingWithSchemaEvolution) { @@ -2419,6 +2452,27 @@ TEST_P(WriteAndReadInteTest, TestSharedShreddingWithSchemaEvolution) { [0, [["a", 32]], [["c", 52]], "new-2"] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr f0_access_field, + BuildMapSharedShreddingAccessField(fields_v1[0], {"z", "a"})); + ASSERT_OK_AND_ASSIGN(std::shared_ptr f2_access_field, + BuildMapSharedShreddingAccessField(fields_v1[3], {"x", "c"})); + auto read_schema = + arrow::schema({f0_access_field, f2_access_field, arrow::field("k2", arrow::utf8())}); + expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + f0_access_field, + f2_access_field, + arrow::field("k2", arrow::utf8()), + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options_v1, read_schema, expected_type, + R"([ + [0, [11, 10], null, "old-1"], + [0, [null, 12], null, "old-2"], + [0, [31, 30], [51, 50], "new-1"], + [0, [null, 32], [null, 52], "new-2"] + ])")); + ASSERT_TRUE(success); } // Verify storage-layout evolution: default->shared-shredding. @@ -2480,6 +2534,23 @@ TEST_P(WriteAndReadInteTest, TestMapStorageLayoutDefaultToSharedShredding) { [0, 4, [["a", 40]]] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"z", "a"})); + auto read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options_v1, read_schema, expected_type, + R"([ + [0, 1, [11, 10]], + [0, 2, null], + [0, 3, [31, 30]], + [0, 4, [null, 40]] + ])")); + ASSERT_TRUE(success); } // Verify storage-layout evolution: shared-shredding->default. @@ -2704,6 +2775,22 @@ TEST_P(WriteAndReadInteTest, TestSharedShreddingWithStructValue) { [0, 3, null] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"a", "z"})); + auto read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options, read_schema, expected_type, + R"([ + [0, 1, [["alice", 10], ["zoe", 11]]], + [0, 2, [["amy", null], null]], + [0, 3, null] + ])")); + ASSERT_TRUE(success); } TEST_P(WriteAndReadInteTest, TestMapSharedShreddingWithComplexValue) { @@ -2798,6 +2885,26 @@ TEST_P(WriteAndReadInteTest, TestMapSharedShreddingWithComplexValue) { [0, 3, null] ])")); ASSERT_TRUE(selected_success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"z", "a"})); + read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(bool access_success, + ReadAndCheckWithReadSchema(options, read_schema, expected_type, + R"([ + [0, 1, [ + ["zeta", [9], [["iz", 90]]], + ["alpha", [1, 2], [["ia", 10], ["ib", 20]]] + ]], + [0, 2, [null, ["amy", null, [["ia", 30]]]]], + [0, 3, null] + ])")); + ASSERT_TRUE(access_success); } TEST_P(WriteAndReadInteTest, TestMapSharedShreddingWithAllSupportedComplexValueTypes) { @@ -2926,6 +3033,28 @@ TEST_P(WriteAndReadInteTest, TestMapSharedShreddingWithAllSupportedComplexValueT [0, 3, []] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"fixed-a"})); + auto read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options, read_schema, expected_type, + R"([ + [0, 1, [[ + true, 1, 2, 3, 4, 5.5, 6.25, "str", "bin", + "12345678.90", "123456789012345678.12345", 19500, + "2023-11-14 22:13:20.123", "2023-11-14 22:13:20.123456789", + "2023-11-14 22:13:20.123", "2023-11-14 22:13:20.123456", + [7, null, 8], [["m1", 10], ["m2", null]], ["nested", 99] + ]]], + [0, 2, null], + [0, 3, [null]] + ])")); + ASSERT_TRUE(success); } TEST_P(WriteAndReadInteTest, TestMapSharedShreddingStructValueSchemaEvolutionReadFails) { @@ -3122,6 +3251,23 @@ TEST_P(WriteAndReadInteTest, TestOrcDictionaryLazyDecodingWithSharedShredding) { [0, 4, [["a", "red"]]] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"z", "a"})); + auto read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options_v1, read_schema, expected_type, + R"([ + [0, 1, ["blue", "red"]], + [0, 2, ["green", "red"]], + [0, 3, ["yellow", "red"]], + [0, 4, [null, "red"]] + ])")); + ASSERT_TRUE(success); } // Verify shared-shredding in the PK read path. @@ -3177,6 +3323,22 @@ TEST_P(WriteAndReadInteTest, TestPkSharedShreddingMap) { [0, 3, [["c", 30]]] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"a", "z"})); + auto read_schema = arrow::schema({arrow::field("pk", arrow::int32()), access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("pk", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options, read_schema, expected_type, + R"([ + [0, 1, [100, 101]], + [0, 2, [null, null]], + [0, 3, [null, null]] + ])")); + ASSERT_TRUE(success); } TEST_P(WriteAndReadInteTest, TestSharedShreddingPartialKeyRecallWithOverflow) { @@ -3290,6 +3452,28 @@ TEST_P(WriteAndReadInteTest, TestSharedShreddingPartialKeyRecallWithOverflow) { ])")); ASSERT_TRUE(success); } + + // Sub-case 4: selected keys are exposed as STRUCT children instead of a filtered MAP. + { + ASSERT_OK_AND_ASSIGN( + std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(arrow::field("tags", map_type), {"c", "a"})); + + auto read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(bool success, + ReadAndCheckWithReadSchema(options, read_schema, expected_type, + R"([ + [0, 1, [3, 1]], + [0, 2, [null, 10]], + [0, 3, null] + ])")); + ASSERT_TRUE(success); + } } TEST_P(WriteAndReadInteTest, TestSharedShreddingPartialKeyRecallWithNullOrMissingKey) { @@ -3381,6 +3565,26 @@ TEST_P(WriteAndReadInteTest, TestSharedShreddingPartialKeyRecallWithNullOrMissin ])")); ASSERT_TRUE(success); } + + // Sub-case 5: expose an existing and a never-written key as STRUCT children. + { + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"a", "nonexistent"})); + auto read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(bool success, + ReadAndCheckWithReadSchema(options, read_schema, expected_type, + R"([ + [0, 1, null], + [0, 2, [null, null]], + [0, 3, [30, null]] + ])")); + ASSERT_TRUE(success); + } } TEST_P(WriteAndReadInteTest, TestSharedShreddingPartialKeyRecallMultipleColumns) { @@ -3470,6 +3674,29 @@ TEST_P(WriteAndReadInteTest, TestSharedShreddingPartialKeyRecallMultipleColumns) ])")); ASSERT_TRUE(success); } + + // Sub-case 3: expose selected keys from multiple MAP columns as independent STRUCTs. + { + ASSERT_OK_AND_ASSIGN(std::shared_ptr tags_access_field, + BuildMapSharedShreddingAccessField(fields[1], {"a", "b"})); + ASSERT_OK_AND_ASSIGN(std::shared_ptr metrics_access_field, + BuildMapSharedShreddingAccessField(fields[2], {"x"})); + auto read_schema = arrow::schema( + {arrow::field("id", arrow::int32()), tags_access_field, metrics_access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + tags_access_field, + metrics_access_field, + }); + ASSERT_OK_AND_ASSIGN(bool success, + ReadAndCheckWithReadSchema(options, read_schema, expected_type, + R"([ + [0, 1, [1, 2], [100]], + [0, 2, [10, null], [1000]] + ])")); + ASSERT_TRUE(success); + } } TEST_P(WriteAndReadInteTest, TestMapStorageLayoutDefaultToSharedShreddingPartialKeyRecall) { @@ -3543,6 +3770,23 @@ TEST_P(WriteAndReadInteTest, TestMapStorageLayoutDefaultToSharedShreddingPartial [0, 4, []] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(arrow::field("tags", map_type), {"a"})); + read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options_v1, read_schema, expected_type, + R"([ + [0, 1, [10]], + [0, 2, null], + [0, 3, [30]], + [0, 4, [null]] + ])")); + ASSERT_TRUE(success); } TEST_P(WriteAndReadInteTest, TestMapStorageLayoutSharedShreddingToDefaultPartialKeyRecall) { @@ -3620,6 +3864,24 @@ TEST_P(WriteAndReadInteTest, TestMapStorageLayoutSharedShreddingToDefaultPartial [0, 4, null] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"a"})); + read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_NOK_WITH_MSG(ReadAndCheckWithReadSchema(options_v1, read_schema, expected_type, + R"([ + [0, 1, [10]], + [0, 2, [null]], + [0, 3, [30]], + [0, 4, null] + ])"), + "Selected-key MAP pushdown only supports top-level shared-shredding MAP " + "field: tags"); } TEST_P(WriteAndReadInteTest, TestSharedShreddingDuplicateSelectedKeys) { @@ -3720,6 +3982,22 @@ TEST_P(WriteAndReadInteTest, TestSharedShreddingAllNullMapColumn) { [0, 3, null] ])")); ASSERT_TRUE(success); + + ASSERT_OK_AND_ASSIGN(std::shared_ptr access_field, + BuildMapSharedShreddingAccessField(fields[1], {"a"})); + auto read_schema = arrow::schema({arrow::field("id", arrow::int32()), access_field}); + auto expected_type = arrow::struct_({ + arrow::field("_VALUE_KIND", arrow::int8()), + arrow::field("id", arrow::int32()), + access_field, + }); + ASSERT_OK_AND_ASSIGN(success, ReadAndCheckWithReadSchema(options, read_schema, expected_type, + R"([ + [0, 1, null], + [0, 2, null], + [0, 3, null] + ])")); + ASSERT_TRUE(success); } INSTANTIATE_TEST_SUITE_P(FileFormatAndFileSystem, WriteAndReadInteTest,