Skip to content

Commit aa9fbf5

Browse files
fix: nested list schema evolution (alibaba#446)
1 parent 6366ca6 commit aa9fbf5

8 files changed

Lines changed: 673 additions & 38 deletions

src/paimon/common/reader/data_evolution_file_reader.cpp

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@
1616

1717
#include "paimon/common/reader/data_evolution_file_reader.h"
1818

19+
#include "arrow/array/array_nested.h"
20+
#include "arrow/array/util.h"
1921
#include "arrow/c/abi.h"
2022
#include "arrow/c/bridge.h"
2123
#include "fmt/format.h"
@@ -25,6 +27,7 @@
2527
#include "paimon/common/utils/arrow/status_utils.h"
2628

2729
namespace paimon {
30+
2831
Result<std::unique_ptr<DataEvolutionFileReader>> DataEvolutionFileReader::Create(
2932
std::vector<std::unique_ptr<BatchReader>>&& readers,
3033
const std::shared_ptr<arrow::Schema>& read_schema, int32_t read_batch_size,
@@ -82,6 +85,7 @@ Result<BatchReader::ReadBatchWithBitmap> DataEvolutionFileReader::NextBatchWithB
8285
}
8386
const auto& sub_array = array_for_each_reader[reader_offsets_[i]];
8487
assert(sub_array->num_fields() > field_offsets_[i]);
88+
// Each file is already aligned to its read schema by its FieldMappingReader.
8589
target_sub_array_vec.push_back(sub_array->field(field_offsets_[i]));
8690
}
8791
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(

src/paimon/core/io/field_mapping_reader.cpp

Lines changed: 20 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,11 @@ Result<std::unique_ptr<FieldMappingReader>> FieldMappingReader::Create(
149149
if (mapping_reader->non_partition_info_.cast_executors[i] != nullptr) {
150150
mapping_reader->need_casting_ = true;
151151
}
152+
// A differing nested type needs the AlignArrayToReadType reshape below.
153+
if (!mapping_reader->non_partition_info_.non_partition_data_schema[i].Type()->Equals(
154+
*mapping_reader->non_partition_info_.non_partition_read_schema[i].Type())) {
155+
mapping_reader->need_casting_ = true;
156+
}
152157
// Field name change (RENAME COLUMN) also requires mapping: data schema
153158
// carries the file's physical name while read schema carries the
154159
// post-rename logical name. If we skipped mapping, the inner reader's
@@ -201,6 +206,7 @@ Result<std::shared_ptr<arrow::Array>> FieldMappingReader::CastNonPartitionArrayI
201206
casted_array.reserve(field_count);
202207
casted_field_names.reserve(field_count);
203208
for (int32_t i = 0; i < field_count; i++) {
209+
std::shared_ptr<arrow::Array> column;
204210
if (non_partition_info_.cast_executors[i] != nullptr) {
205211
auto single_column_array = struct_array->field(i);
206212
// if src array is dict, cast to string first
@@ -213,18 +219,27 @@ Result<std::shared_ptr<arrow::Array>> FieldMappingReader::CastNonPartitionArrayI
213219
arrow::compute::CastOptions::Safe(), arrow_pool_.get()));
214220
}
215221
PAIMON_ASSIGN_OR_RAISE(
216-
std::shared_ptr<arrow::Array> casted,
222+
column,
217223
non_partition_info_.cast_executors[i]->Cast(
218224
single_column_array, non_partition_info_.non_partition_read_schema[i].Type(),
219225
arrow_pool_.get()));
220-
casted_array.push_back(casted);
221-
casted_field_names.push_back(non_partition_info_.non_partition_data_schema[i].Name());
222226
} else {
223227
// read and data type may both be string type, but after adapter transform, type may be
224228
// dictionary, need reconstruct struct type
225-
casted_array.push_back(struct_array->field(i));
226-
casted_field_names.push_back(non_partition_info_.non_partition_data_schema[i].Name());
229+
column = struct_array->field(i);
230+
}
231+
// Null-fill nested fields added by schema evolution. Only when the data and
232+
// read types differ -- the reader may hand back a dictionary-encoded array
233+
// for an unchanged type, which is not a reshape target.
234+
if (!non_partition_info_.non_partition_data_schema[i].Type()->Equals(
235+
*non_partition_info_.non_partition_read_schema[i].Type())) {
236+
PAIMON_ASSIGN_OR_RAISE(
237+
column, NestedProjectionUtils::AlignArrayToReadType(
238+
column, non_partition_info_.non_partition_read_schema[i].Type(),
239+
arrow_pool_.get()));
227240
}
241+
casted_array.push_back(column);
242+
casted_field_names.push_back(non_partition_info_.non_partition_data_schema[i].Name());
228243
}
229244
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Array> arrow_array,
230245
arrow::StructArray::Make(casted_array, casted_field_names));

src/paimon/core/io/field_mapping_reader_test.cpp

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -448,6 +448,91 @@ TEST_F(FieldMappingReaderTest, TestDictionaryTypeWithSchemaEvolution) {
448448
partition, expected_array);
449449
}
450450

451+
TEST_F(FieldMappingReaderTest, TestSchemaEvolutionAddedFieldInsideList) {
452+
// A field `c`(id=12) was added inside the list's struct after the file was
453+
// written. Reading the old file with the new schema must null-fill `c`.
454+
auto id_field = [](const std::string& name, const std::shared_ptr<arrow::DataType>& type,
455+
int32_t id) {
456+
return DataField::ConvertDataFieldToArrowField(DataField(id, arrow::field(name, type)));
457+
};
458+
auto data_struct =
459+
arrow::struct_({id_field("a", arrow::int32(), 10), id_field("b", arrow::utf8(), 11)});
460+
auto read_struct =
461+
arrow::struct_({id_field("a", arrow::int32(), 10), id_field("b", arrow::utf8(), 11),
462+
id_field("c", arrow::int32(), 12)});
463+
std::vector<DataField> data_fields = {
464+
DataField(100, arrow::field("items", arrow::list(arrow::field("item", data_struct))))};
465+
std::vector<DataField> read_fields = {
466+
DataField(100, arrow::field("items", arrow::list(arrow::field("item", read_struct))))};
467+
auto data_schema = DataField::ConvertDataFieldsToArrowSchema(data_fields);
468+
auto read_schema = DataField::ConvertDataFieldsToArrowSchema(read_fields);
469+
470+
auto data_array = std::dynamic_pointer_cast<arrow::StructArray>(
471+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(data_schema->fields()), R"([
472+
[[[1, "x"], [2, "y"]]],
473+
[[[3, "z"]]]
474+
])")
475+
.ValueOrDie());
476+
477+
ASSERT_OK_AND_ASSIGN(auto mapping_builder,
478+
FieldMappingBuilder::Create(read_schema, /*partition_keys=*/{},
479+
/*predicate=*/nullptr));
480+
ASSERT_OK_AND_ASSIGN(auto mapping, mapping_builder->CreateFieldMapping(data_fields));
481+
auto mock = std::make_unique<MockFileBatchReader>(
482+
data_array, arrow::struct_(data_schema->fields()), /*read_batch_size=*/8);
483+
ASSERT_OK_AND_ASSIGN(auto reader, FieldMappingReader::Create(
484+
read_schema->num_fields(), std::move(mock),
485+
BinaryRow::EmptyRow(), std::move(mapping),
486+
/*skip_map_selected_keys_filter_field_ids=*/{}, pool_));
487+
ASSERT_OK_AND_ASSIGN(auto result_array, ReadResultCollector::CollectResult(reader.get()));
488+
489+
auto expect_array =
490+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(read_schema->fields()), R"([
491+
[[[1, "x", null], [2, "y", null]]],
492+
[[[3, "z", null]]]
493+
])")
494+
.ValueOrDie();
495+
auto expected_chunk = std::make_shared<arrow::ChunkedArray>(arrow::ArrayVector({expect_array}));
496+
ASSERT_TRUE(result_array->type()->Equals(expected_chunk->type()))
497+
<< result_array->type()->ToString() << " vs " << expected_chunk->type()->ToString();
498+
ASSERT_TRUE(result_array->Equals(expected_chunk))
499+
<< result_array->ToString() << " vs " << expected_chunk->ToString();
500+
}
501+
502+
TEST_F(FieldMappingReaderTest, TestSchemaEvolutionAddedFieldInsideListOrc) {
503+
// ORC round-trip: added field inside a list's struct is null-filled.
504+
auto id_field = [](const std::string& name, const std::shared_ptr<arrow::DataType>& type,
505+
int32_t id) {
506+
return DataField::ConvertDataFieldToArrowField(DataField(id, arrow::field(name, type)));
507+
};
508+
auto data_struct =
509+
arrow::struct_({id_field("a", arrow::int32(), 10), id_field("b", arrow::int32(), 11)});
510+
auto read_struct =
511+
arrow::struct_({id_field("a", arrow::int32(), 10), id_field("b", arrow::int32(), 11),
512+
id_field("c", arrow::int32(), 12)});
513+
std::vector<DataField> data_fields = {
514+
DataField(100, arrow::field("items", arrow::list(arrow::field("item", data_struct))))};
515+
std::vector<DataField> read_fields = {
516+
DataField(100, arrow::field("items", arrow::list(arrow::field("item", read_struct))))};
517+
auto data_schema = DataField::ConvertDataFieldsToArrowSchema(data_fields);
518+
auto read_schema = DataField::ConvertDataFieldsToArrowSchema(read_fields);
519+
520+
auto data_array = std::dynamic_pointer_cast<arrow::StructArray>(
521+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(data_schema->fields()), R"([
522+
[[[1, 2], [3, 4]]],
523+
[[[5, 6]]]
524+
])")
525+
.ValueOrDie());
526+
auto expect_array =
527+
arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(read_schema->fields()), R"([
528+
[[[1, 2, null], [3, 4, null]]],
529+
[[[5, 6, null]]]
530+
])")
531+
.ValueOrDie();
532+
CheckResult(data_schema, data_array, read_schema, /*predicate=*/nullptr, /*partition_keys=*/{},
533+
BinaryRow::EmptyRow(), expect_array);
534+
}
535+
451536
TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithModifyType) {
452537
std::vector<DataField> data_fields = {DataField(0, arrow::field("f0", arrow::utf8())),
453538
DataField(1, arrow::field("f1", arrow::float32())),

src/paimon/core/utils/field_mapping.cpp

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -172,10 +172,11 @@ Result<std::vector<std::shared_ptr<CastExecutor>>> FieldMappingBuilder::CreateDa
172172
FieldTypeUtils::ConvertToFieldType(data_fields[i].Type()->id()));
173173

174174
if (!read_fields[i].Type()->Equals(data_fields[i].Type())) {
175-
if (read_type == FieldType::STRUCT) {
176-
// STRUCT may still differ by nested pruning shape. No cast is
177-
// needed — type pruning is handled by PruneDataType during
178-
// field mapping construction.
175+
auto read_type_id = read_fields[i].Type()->id();
176+
if (read_type_id == arrow::Type::STRUCT || read_type_id == arrow::Type::LIST ||
177+
read_type_id == arrow::Type::MAP) {
178+
// Nested type differs by pruning/evolution; the reader's reshape
179+
// handles it, no scalar cast.
179180
cast_executors.push_back(nullptr);
180181
continue;
181182
}

0 commit comments

Comments
 (0)