Skip to content

Commit 5cd814d

Browse files
authored
fix: fix compaction crash when PK fields are not at the beginning of table schema (alibaba#202)
1 parent 6ad44ae commit 5cd814d

3 files changed

Lines changed: 720 additions & 13 deletions

File tree

src/paimon/core/operation/merge_file_split_read.cpp

Lines changed: 24 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -285,6 +285,24 @@ Status MergeFileSplitRead::GenerateKeyValueReadSchema(
285285
std::shared_ptr<arrow::Schema>* value_schema, std::shared_ptr<arrow::Schema>* read_schema,
286286
std::shared_ptr<FieldsComparator>* key_comparator,
287287
std::shared_ptr<FieldsComparator>* sequence_fields_comparator) {
288+
PAIMON_ASSIGN_OR_RAISE(std::vector<DataField> trimmed_key_fields,
289+
table_schema.TrimmedPrimaryKeyFields());
290+
PAIMON_ASSIGN_OR_RAISE(*key_comparator, FieldsComparator::Create(trimmed_key_fields,
291+
/*is_ascending_order=*/true));
292+
const auto& table_fields = table_schema.Fields();
293+
auto table_fields_schema = DataField::ConvertDataFieldsToArrowSchema(table_fields);
294+
if (table_fields_schema->Equals(raw_read_schema)) {
295+
// Short-circuit: if raw_read_schema is the same as the table schema,
296+
// use the table schema field order directly (for compact process).
297+
*value_schema = table_fields_schema;
298+
// sequence_fields_comparator
299+
PAIMON_ASSIGN_OR_RAISE(
300+
*sequence_fields_comparator,
301+
PrimaryKeyTableUtils::CreateSequenceFieldsComparator(table_fields, options));
302+
*read_schema = SpecialFields::CompleteSequenceAndValueKindField(*value_schema);
303+
return Status::OK();
304+
}
305+
288306
// 1. add user raw read schema to need_fields
289307
PAIMON_ASSIGN_OR_RAISE(std::vector<DataField> need_fields,
290308
DataField::ConvertArrowSchemaToDataFields(raw_read_schema));
@@ -302,10 +320,10 @@ Status MergeFileSplitRead::GenerateKeyValueReadSchema(
302320
// 3. split need_fields to key and non-key fields
303321
std::vector<DataField> key_fields;
304322
std::vector<DataField> non_key_fields;
305-
PAIMON_ASSIGN_OR_RAISE(std::vector<std::string> trimmed_key_fields,
323+
PAIMON_ASSIGN_OR_RAISE(std::vector<std::string> trimmed_key_names,
306324
table_schema.TrimmedPrimaryKeys());
307325
PAIMON_RETURN_NOT_OK(
308-
SplitKeyAndNonKeyField(trimmed_key_fields, need_fields, &key_fields, &non_key_fields));
326+
SplitKeyAndNonKeyField(trimmed_key_names, need_fields, &key_fields, &non_key_fields));
309327

310328
// 4. construct value fields: key fields are put before non-key fields
311329
std::vector<DataField> value_fields;
@@ -316,17 +334,10 @@ Status MergeFileSplitRead::GenerateKeyValueReadSchema(
316334
PAIMON_ASSIGN_OR_RAISE(
317335
*sequence_fields_comparator,
318336
PrimaryKeyTableUtils::CreateSequenceFieldsComparator(value_fields, options));
319-
// 6. complete key fields to all trimmed primary key
320-
key_fields.clear();
321-
PAIMON_ASSIGN_OR_RAISE(key_fields, table_schema.GetFields(trimmed_key_fields));
322-
PAIMON_ASSIGN_OR_RAISE(*key_comparator, FieldsComparator::Create(key_fields,
323-
/*is_ascending_order=*/true));
324-
// 7. construct actual read fields: special + key + non-key value
325-
std::vector<DataField> read_fields;
326-
std::vector<DataField> special_fields(
327-
{SpecialFields::SequenceNumber(), SpecialFields::ValueKind()});
328-
read_fields.insert(read_fields.end(), special_fields.begin(), special_fields.end());
329-
read_fields.insert(read_fields.end(), key_fields.begin(), key_fields.end());
337+
// 6. construct actual read fields: special + key + non-key value
338+
std::vector<DataField> read_fields = {SpecialFields::SequenceNumber(),
339+
SpecialFields::ValueKind()};
340+
read_fields.insert(read_fields.end(), trimmed_key_fields.begin(), trimmed_key_fields.end());
330341
read_fields.insert(read_fields.end(), non_key_fields.begin(), non_key_fields.end());
331342
*read_schema = DataField::ConvertDataFieldsToArrowSchema(read_fields);
332343
return Status::OK();

test/inte/CMakeLists.txt

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -97,4 +97,11 @@ if(PAIMON_BUILD_TESTS)
9797
test_utils_static
9898
${GTEST_LINK_TOOLCHAIN})
9999

100+
add_paimon_test(pk_compaction_inte_test
101+
STATIC_LINK_LIBS
102+
paimon_shared
103+
${TEST_STATIC_LINK_LIBS}
104+
test_utils_static
105+
${GTEST_LINK_TOOLCHAIN})
106+
100107
endif()

0 commit comments

Comments
 (0)