Skip to content

Commit afc8471

Browse files
authored
fix: optimize shared-shredding read & fix ORC read-size estimation for nested columns (#216)
1 parent c7f23c6 commit afc8471

4 files changed

Lines changed: 413 additions & 82 deletions

File tree

src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp

Lines changed: 154 additions & 78 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,6 @@
3232
#include "paimon/common/utils/arrow/mem_utils.h"
3333
#include "paimon/common/utils/arrow/status_utils.h"
3434
#include "paimon/common/utils/checked_cast.h"
35-
#include "paimon/core/casting/casting_utils.h"
3635
#include "paimon/core/utils/nested_projection_utils.h"
3736

3837
namespace paimon {
@@ -110,6 +109,56 @@ class SharedSelectedKeysReadPlan : public MapFieldReadPlan {
110109
std::vector<SelectedKey> selected_keys_;
111110
};
112111

112+
Result<std::shared_ptr<arrow::Array>> MaskSinglePhysicalColumn(
113+
const std::shared_ptr<arrow::StructArray>& physical_struct_array,
114+
const std::shared_ptr<arrow::ListArray>& field_mapping_array,
115+
const std::shared_ptr<arrow::Int32Array>& field_mapping_values,
116+
const std::shared_ptr<arrow::Array>& physical_column_array, int32_t physical_column_id,
117+
int32_t field_id, const std::string& field_name, arrow::MemoryPool* arrow_pool) {
118+
int64_t row_count = physical_struct_array->length();
119+
if (physical_column_array->length() != row_count) {
120+
return Status::Invalid("shared-shredding physical column length does not match row count");
121+
}
122+
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Buffer> validity,
123+
arrow::AllocateEmptyBitmap(row_count, arrow_pool));
124+
int64_t valid_count = 0;
125+
for (int64_t row = 0; row < row_count; ++row) {
126+
if (physical_struct_array->IsNull(row)) {
127+
continue;
128+
}
129+
if (field_mapping_array->IsNull(row)) {
130+
return Status::Invalid(fmt::format(
131+
"__field_mapping cannot be null in non-null shared-shredding row for field {}",
132+
field_name));
133+
}
134+
int32_t mapping_offset = field_mapping_array->value_offset(row);
135+
int32_t mapping_length = field_mapping_array->value_length(row);
136+
if (physical_column_id < 0 || physical_column_id >= mapping_length) {
137+
return Status::Invalid("physical column id is out of __field_mapping range");
138+
}
139+
int32_t mapping_index = mapping_offset + physical_column_id;
140+
if (field_mapping_values->IsNull(mapping_index)) {
141+
return Status::Invalid("__field_mapping element cannot be null");
142+
}
143+
if (field_mapping_values->Value(mapping_index) != field_id ||
144+
physical_column_array->IsNull(row)) {
145+
continue;
146+
}
147+
arrow::bit_util::SetBit(validity->mutable_data(), row);
148+
++valid_count;
149+
}
150+
151+
// Replace only the top-level validity; offsets, values, and nested children stay shared.
152+
std::shared_ptr<arrow::ArrayData> result_data = physical_column_array->data()->Copy();
153+
if (result_data->buffers.empty()) {
154+
return Status::Invalid("shared-shredding physical column has no validity buffer slot");
155+
}
156+
int64_t null_count = row_count - valid_count;
157+
result_data->buffers[0] = null_count == 0 ? nullptr : std::move(validity);
158+
result_data->SetNullCount(null_count);
159+
return arrow::MakeArray(std::move(result_data));
160+
}
161+
113162
class DefaultSelectedKeysReadPlan : public MapFieldReadPlan {
114163
public:
115164
DefaultSelectedKeysReadPlan(const std::shared_ptr<arrow::Field>& logical_field,
@@ -386,12 +435,10 @@ Result<std::shared_ptr<arrow::Array>> FullMapReadPlan::Materialize(
386435
std::shared_ptr<arrow::MapArray> overflow_array;
387436
CollectPhysicalColumns(physical_struct_array, &physical_column_name_to_array, &overflow_array);
388437
for (auto& [_, physical_column_array] : physical_column_name_to_array) {
389-
if (physical_column_array->type_id() == arrow::Type::DICTIONARY) {
390-
PAIMON_ASSIGN_OR_RAISE(
391-
physical_column_array,
392-
CastingUtils::Cast(physical_column_array, logical_map_type_->item_type(),
393-
arrow::compute::CastOptions::Safe(), arrow_pool));
394-
}
438+
PAIMON_ASSIGN_OR_RAISE(
439+
physical_column_array,
440+
NestedProjectionUtils::AlignArrayToReadType(
441+
physical_column_array, logical_map_type_->item_type(), arrow_pool));
395442
}
396443

397444
std::shared_ptr<arrow::Int32Array> overflow_keys;
@@ -406,12 +453,9 @@ Result<std::shared_ptr<arrow::Array>> FullMapReadPlan::Materialize(
406453
if (!overflow_items) {
407454
return Status::Invalid("__overflow map item array is null");
408455
}
409-
if (overflow_items->type_id() == arrow::Type::DICTIONARY) {
410-
PAIMON_ASSIGN_OR_RAISE(
411-
overflow_items,
412-
CastingUtils::Cast(overflow_items, logical_map_type_->item_type(),
413-
arrow::compute::CastOptions::Safe(), arrow_pool));
414-
}
456+
PAIMON_ASSIGN_OR_RAISE(overflow_items,
457+
NestedProjectionUtils::AlignArrayToReadType(
458+
overflow_items, logical_map_type_->item_type(), arrow_pool));
415459
}
416460

417461
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::unique_ptr<arrow::ArrayBuilder> map_builder_base,
@@ -530,12 +574,9 @@ Result<std::shared_ptr<arrow::Array>> SharedSelectedKeysReadPlan::Materialize(
530574
std::shared_ptr<arrow::MapArray> overflow_array;
531575
CollectPhysicalColumns(physical_struct_array, &physical_column_name_to_array, &overflow_array);
532576
for (auto& [_, physical_column_array] : physical_column_name_to_array) {
533-
if (physical_column_array->type_id() == arrow::Type::DICTIONARY) {
534-
PAIMON_ASSIGN_OR_RAISE(
535-
physical_column_array,
536-
CastingUtils::Cast(physical_column_array, value_type,
537-
arrow::compute::CastOptions::Safe(), arrow_pool));
538-
}
577+
PAIMON_ASSIGN_OR_RAISE(physical_column_array,
578+
NestedProjectionUtils::AlignArrayToReadType(physical_column_array,
579+
value_type, arrow_pool));
539580
}
540581

541582
std::shared_ptr<arrow::Int32Array> overflow_keys;
@@ -550,65 +591,89 @@ Result<std::shared_ptr<arrow::Array>> SharedSelectedKeysReadPlan::Materialize(
550591
if (!overflow_items) {
551592
return Status::Invalid("__overflow map item array is null");
552593
}
553-
if (overflow_items->type_id() == arrow::Type::DICTIONARY) {
554-
PAIMON_ASSIGN_OR_RAISE(
555-
overflow_items,
556-
CastingUtils::Cast(overflow_items, value_type, arrow::compute::CastOptions::Safe(),
557-
arrow_pool));
558-
}
594+
PAIMON_ASSIGN_OR_RAISE(overflow_items, NestedProjectionUtils::AlignArrayToReadType(
595+
overflow_items, value_type, arrow_pool));
559596
}
560597

561-
std::unique_ptr<arrow::ArrayBuilder> access_builder_base;
562-
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(access_builder_base,
563-
arrow::MakeBuilder(LogicalField()->type(), arrow_pool));
564-
if (!access_builder_base || !access_builder_base->type() ||
565-
access_builder_base->type()->id() != arrow::Type::STRUCT) {
566-
return Status::Invalid(
567-
fmt::format("selected-key MAP field {} is not a STRUCT", LogicalField()->name()));
568-
}
569-
auto* access_builder = checked_cast<arrow::StructBuilder*>(access_builder_base.get());
570-
PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Reserve(physical_struct_array->length()));
571-
572-
for (int64_t row = 0; row < physical_struct_array->length(); ++row) {
573-
if (physical_struct_array->IsNull(row)) {
574-
PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->AppendNull());
598+
int64_t row_count = physical_struct_array->length();
599+
arrow::ArrayVector selected_key_arrays;
600+
selected_key_arrays.reserve(selected_keys_.size());
601+
for (int32_t key_index = 0; key_index < selected_keys_type->num_fields(); ++key_index) {
602+
const SelectedKey& selected_key = selected_keys_[key_index];
603+
if (selected_key.field_id < 0) {
604+
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(
605+
std::shared_ptr<arrow::Array> null_array,
606+
arrow::MakeArrayOfNull(selected_keys_type->field(key_index)->type(), row_count,
607+
arrow_pool));
608+
selected_key_arrays.push_back(std::move(null_array));
575609
continue;
576610
}
577-
if (field_mapping_array->IsNull(row)) {
578-
return Status::Invalid(fmt::format(
579-
"__field_mapping cannot be null in non-null shared-shredding row for field {}",
580-
LogicalField()->name()));
611+
612+
if (selected_key.candidate_columns.size() == 1 && !selected_key.may_use_overflow) {
613+
int32_t physical_column_id = selected_key.candidate_columns[0];
614+
std::string physical_column_name =
615+
MapSharedShreddingDefine::PhysicalColumnName(physical_column_id);
616+
auto physical_column_iter = physical_column_name_to_array.find(physical_column_name);
617+
if (physical_column_iter == physical_column_name_to_array.end()) {
618+
return Status::Invalid(
619+
fmt::format("cannot find selected physical column {} for field {}",
620+
physical_column_name, LogicalField()->name()));
621+
}
622+
const std::shared_ptr<arrow::Array>& physical_column_array =
623+
physical_column_iter->second;
624+
if (physical_column_array->offset() == 0) {
625+
PAIMON_ASSIGN_OR_RAISE(
626+
std::shared_ptr<arrow::Array> masked_array,
627+
MaskSinglePhysicalColumn(physical_struct_array, field_mapping_array,
628+
field_mapping_values, physical_column_array,
629+
physical_column_id, selected_key.field_id,
630+
LogicalField()->name(), arrow_pool));
631+
selected_key_arrays.push_back(std::move(masked_array));
632+
continue;
633+
} else {
634+
return Status::Invalid("paimon only supports arrays with zero offset");
635+
}
581636
}
582-
int32_t mapping_offset = field_mapping_array->value_offset(row);
583637

584-
PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Append());
585-
for (int32_t key_index = 0; key_index < selected_keys_type->num_fields(); ++key_index) {
586-
arrow::ArrayBuilder* value_builder = access_builder->field_builder(key_index);
587-
const SelectedKey& selected_key = selected_keys_[key_index];
638+
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(
639+
std::unique_ptr<arrow::ArrayBuilder> value_builder,
640+
arrow::MakeBuilder(selected_keys_type->field(key_index)->type(), arrow_pool));
641+
PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->Reserve(row_count));
642+
for (int64_t row = 0; row < row_count; ++row) {
643+
if (physical_struct_array->IsNull(row)) {
644+
PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->AppendNull());
645+
continue;
646+
}
647+
if (field_mapping_array->IsNull(row)) {
648+
return Status::Invalid(fmt::format(
649+
"__field_mapping cannot be null in non-null shared-shredding row for field {}",
650+
LogicalField()->name()));
651+
}
652+
int32_t mapping_offset = field_mapping_array->value_offset(row);
588653
bool appended = false;
589-
if (selected_key.field_id >= 0) {
590-
for (int32_t physical_column_id : selected_key.candidate_columns) {
591-
int32_t mapping_index = mapping_offset + physical_column_id;
592-
if (field_mapping_values->IsNull(mapping_index)) {
593-
return Status::Invalid("__field_mapping element cannot be null");
594-
}
595-
if (field_mapping_values->Value(mapping_index) != selected_key.field_id) {
596-
continue;
597-
}
598-
std::string physical_column_name =
599-
MapSharedShreddingDefine::PhysicalColumnName(physical_column_id);
600-
auto physical_column_iter =
601-
physical_column_name_to_array.find(physical_column_name);
602-
if (physical_column_iter == physical_column_name_to_array.end()) {
603-
return Status::Invalid(
604-
fmt::format("cannot find selected physical column {} for field {}",
605-
physical_column_name, LogicalField()->name()));
606-
}
607-
PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->AppendArraySlice(
608-
*physical_column_iter->second->data(), row, 1));
609-
appended = true;
610-
break;
654+
for (int32_t physical_column_id : selected_key.candidate_columns) {
655+
int32_t mapping_index = mapping_offset + physical_column_id;
656+
if (field_mapping_values->IsNull(mapping_index)) {
657+
return Status::Invalid("__field_mapping element cannot be null");
658+
}
659+
if (field_mapping_values->Value(mapping_index) != selected_key.field_id) {
660+
continue;
611661
}
662+
std::string physical_column_name =
663+
MapSharedShreddingDefine::PhysicalColumnName(physical_column_id);
664+
auto physical_column_iter =
665+
physical_column_name_to_array.find(physical_column_name);
666+
if (physical_column_iter == physical_column_name_to_array.end()) {
667+
return Status::Invalid(
668+
fmt::format("cannot find selected physical column {} for field {}",
669+
physical_column_name, LogicalField()->name()));
670+
}
671+
const std::shared_ptr<arrow::Array>& physical_column_array =
672+
physical_column_iter->second;
673+
PAIMON_RETURN_NOT_OK_FROM_ARROW(
674+
value_builder->AppendArraySlice(*physical_column_array->data(), row, 1));
675+
appended = true;
676+
break;
612677
}
613678

614679
if (!appended && selected_key.may_use_overflow && overflow_array &&
@@ -630,9 +695,24 @@ Result<std::shared_ptr<arrow::Array>> SharedSelectedKeysReadPlan::Materialize(
630695
PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->AppendNull());
631696
}
632697
}
698+
std::shared_ptr<arrow::Array> selected_key_array;
699+
PAIMON_RETURN_NOT_OK_FROM_ARROW(value_builder->Finish(&selected_key_array));
700+
selected_key_arrays.push_back(std::move(selected_key_array));
701+
}
702+
703+
std::shared_ptr<arrow::Buffer> parent_validity;
704+
int64_t parent_null_count = physical_struct_array->null_count();
705+
if (parent_null_count > 0) {
706+
if (physical_struct_array->offset() == 0) {
707+
parent_validity = physical_struct_array->null_bitmap();
708+
} else {
709+
return Status::Invalid("paimon only supports arrays with zero offset");
710+
}
633711
}
634-
std::shared_ptr<arrow::Array> result;
635-
PAIMON_RETURN_NOT_OK_FROM_ARROW(access_builder->Finish(&result));
712+
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(
713+
std::shared_ptr<arrow::StructArray> result,
714+
arrow::StructArray::Make(selected_key_arrays, selected_keys_type->fields(),
715+
std::move(parent_validity), parent_null_count));
636716
return result;
637717
}
638718

@@ -646,14 +726,10 @@ Result<std::shared_ptr<arrow::Array>> DefaultSelectedKeysReadPlan::Materialize(
646726
}
647727
auto map_array = checked_pointer_cast<arrow::MapArray>(physical_array);
648728
auto selected_keys_type = checked_pointer_cast<arrow::StructType>(LogicalField()->type());
649-
auto physical_map_type = checked_pointer_cast<arrow::MapType>(PhysicalReadField()->type());
650729

651730
std::shared_ptr<arrow::Array> items = map_array->items();
652-
if (items->type_id() == arrow::Type::DICTIONARY) {
653-
PAIMON_ASSIGN_OR_RAISE(items,
654-
CastingUtils::Cast(items, physical_map_type->item_type(),
655-
arrow::compute::CastOptions::Safe(), arrow_pool));
656-
}
731+
PAIMON_ASSIGN_OR_RAISE(items, NestedProjectionUtils::AlignArrayToReadType(
732+
items, selected_keys_type->field(0)->type(), arrow_pool));
657733
std::shared_ptr<arrow::Array> keys = map_array->keys();
658734
std::unique_ptr<arrow::ArrayBuilder> access_builder_base;
659735
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(access_builder_base,

0 commit comments

Comments
 (0)