diff --git a/paimon-python/pypaimon/read/split_read.py b/paimon-python/pypaimon/read/split_read.py index c5300a620fa0..9073072b48fe 100644 --- a/paimon-python/pypaimon/read/split_read.py +++ b/paimon-python/pypaimon/read/split_read.py @@ -632,15 +632,14 @@ def _file_read_fields(self, file: DataFileMeta) -> Optional[List[DataField]]: self._all_data_fields_from(file_schema.fields)) def _create_key_value_fields(self, value_field: List[DataField]): - all_fields: List[DataField] = self.table.fields all_data_fields = [] - for field in all_fields: - if field.name in self.trimmed_primary_key: - key_field_name = f"{KEY_PREFIX}{field.name}" - key_field_id = field.id + KEY_FIELD_ID_START - key_field = DataField(key_field_id, key_field_name, field.type) - all_data_fields.append(key_field) + # Merge keys must follow the same declared order as the sorted files. + for field in self.table.trimmed_primary_keys_fields: + key_field_name = f"{KEY_PREFIX}{field.name}" + key_field_id = field.id + KEY_FIELD_ID_START + key_field = DataField(key_field_id, key_field_name, field.type) + all_data_fields.append(key_field) all_data_fields.append(SpecialFields.SEQUENCE_NUMBER) all_data_fields.append(SpecialFields.VALUE_KIND) diff --git a/paimon-python/pypaimon/tests/reader_primary_key_test.py b/paimon-python/pypaimon/tests/reader_primary_key_test.py index 1ed2148099e5..19a4f1558c31 100644 --- a/paimon-python/pypaimon/tests/reader_primary_key_test.py +++ b/paimon-python/pypaimon/tests/reader_primary_key_test.py @@ -97,6 +97,49 @@ def test_pk_parquet_reader(self): value_kind_field_found, "_VALUE_KIND field should exist in the written parquet file") + def test_composite_primary_key_order(self): + cases = [ + (['customer_id', 'order_id'], []), + (['order_id', 'customer_id'], []), + (['order_id', 'dt', 'customer_id'], ['dt']), + ] + for index, (primary_keys, partition_keys) in enumerate(cases): + with self.subTest(primary_keys=primary_keys): + fields = [('customer_id', pa.int64()), ('order_id', pa.int64()), + ('amount', pa.int64())] + base = {'customer_id': [1, 2, 1, 2], 'order_id': [10, 10, 20, 20], + 'amount': [100, 200, 300, 400]} + update = {'customer_id': [1], 'order_id': [20], 'amount': [350]} + expected = {'customer_id': [1, 1, 2, 2], 'order_id': [10, 20, 10, 20], + 'amount': [100, 350, 200, 400]} + if partition_keys: + fields.append(('dt', pa.string())) + base['dt'] = expected['dt'] = ['p1'] * 4 + update['dt'] = ['p1'] + arrow_schema = pa.schema(fields) + name = 'default.test_composite_pk_order_{}'.format(index) + self.catalog.create_table(name, Schema.from_pyarrow_schema( + arrow_schema, primary_keys=primary_keys, partition_keys=partition_keys, + options={'bucket': '1'}), False) + table = self.catalog.get_table(name) + for data in (base, update): + builder = table.new_batch_write_builder() + writer, committer = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(pa.Table.from_pydict(data, schema=arrow_schema)) + committer.commit(writer.prepare_commit()) + finally: + writer.close() + committer.close() + + actual = self._read_test_table(table.new_read_builder()).sort_by( + [('customer_id', 'ascending'), ('order_id', 'ascending')]) + self.assertEqual(actual.to_pylist(), + pa.Table.from_pydict(expected, schema=arrow_schema).to_pylist()) + projected = self._read_test_table( + table.new_read_builder().with_projection(['amount'])) + self.assertEqual(sorted(projected['amount'].to_pylist()), [100, 200, 350, 400]) + def test_pk_orc_reader(self): schema = Schema.from_pyarrow_schema(self.pa_schema, partition_keys=['dt'], diff --git a/paimon-python/pypaimon/tests/schema_evolution_pk_read_test.py b/paimon-python/pypaimon/tests/schema_evolution_pk_read_test.py index 4df02bb33c46..6fb0afd21d40 100644 --- a/paimon-python/pypaimon/tests/schema_evolution_pk_read_test.py +++ b/paimon-python/pypaimon/tests/schema_evolution_pk_read_test.py @@ -197,6 +197,33 @@ def test_pk_column_position_deduplicate(self): self.assertEqual(self._read_sorted(table), [ {'b': 'b2', 'id': 1, 'a': 'a2'}]) + def test_composite_pk_column_position(self): + s0 = pa.schema([('customer_id', pa.int64()), ('order_id', pa.int64()), + ('amount', pa.int64())]) + table = self._create('composite_pk_position', s0, + primary_keys=('customer_id', 'order_id')) + self._write(table, pa.Table.from_pydict({ + 'customer_id': [1, 1, 2, 2], 'order_id': [10, 20, 10, 20], + 'amount': [100, 200, 300, 400]}, schema=s0)) + + self.catalog.alter_table( + 'default.composite_pk_position', + [SchemaChange.update_column_position(Move.first('order_id'))], False) + table = self.catalog.get_table('default.composite_pk_position') + s1 = pa.schema([('order_id', pa.int64()), ('customer_id', pa.int64()), + ('amount', pa.int64())]) + self._write(table, pa.Table.from_pydict({ + 'order_id': [10], 'customer_id': [2], 'amount': [350]}, schema=s1)) + + builder = table.new_read_builder() + actual = builder.new_read().to_arrow(builder.new_scan().plan().splits()) + self.assertEqual(actual.sort_by([ + ('customer_id', 'ascending'), ('order_id', 'ascending')]).to_pylist(), [ + {'customer_id': 1, 'order_id': 10, 'amount': 100}, + {'customer_id': 1, 'order_id': 20, 'amount': 200}, + {'customer_id': 2, 'order_id': 10, 'amount': 350}, + {'customer_id': 2, 'order_id': 20, 'amount': 400}]) + # -- B7: first-row engine + add column ------------------------------- def test_pk_first_row_add_column(self):