diff --git a/paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java b/paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java index 0cdcb9879106..1d34ad978ce4 100644 --- a/paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java +++ b/paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java @@ -159,7 +159,11 @@ public void before() throws Exception { Files.createDirectories(tempDir.resolve("warehouse")); } - warehouse = new Path(TraceableFileIO.SCHEME + "://" + tempDir.resolve("warehouse")); + warehouse = + new Path( + TraceableFileIO.SCHEME + + "://" + + tempDir.resolve("warehouse").toUri().getPath()); catalog = CatalogFactory.createCatalog(CatalogContext.create(warehouse)); // Create database if it doesn't exist @@ -384,6 +388,37 @@ public void testJavaWriteDynamicBucketHashIndex() throws Exception { } } + @Test + @EnabledIfSystemProperty(named = "run.e2e.tests", matches = "true") + public void testJavaWriteCompositeDatePartition() throws Exception { + for (boolean legacyName : Arrays.asList(true, false)) { + String suffix = legacyName ? "legacy" : "canonical"; + Identifier identifier = identifier("composite_date_java_to_python_" + suffix); + catalog.dropTable(identifier, true); + Schema schema = + Schema.newBuilder() + .column("id", DataTypes.INT()) + .column("day", DataTypes.DATE()) + .column("region", DataTypes.STRING()) + .partitionKeys("day", "region") + .option("partition.legacy-name", Boolean.toString(legacyName)) + .option("bucket", "1") + .option("bucket-key", "id") + .option("write.native.enabled", "false") + .option("commit.native.enabled", "false") + .option("scan.native-plan.enabled", "false") + .build(); + catalog.createTable(identifier, schema, false); + FileStoreTable table = (FileStoreTable) catalog.getTable(identifier); + BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder(); + try (BatchTableWrite write = writeBuilder.newWrite(); + BatchTableCommit commit = writeBuilder.newCommit()) { + write.write(GenericRow.of(1, 1, BinaryString.fromString("a/b"))); + commit.commit(write.prepareCommit()); + } + } + } + @Test @EnabledIfSystemProperty(named = "run.e2e.tests", matches = "true") public void testPKDeletionVectorWrite() throws Exception { @@ -672,6 +707,41 @@ public void testReadPythonDynamicBucketHashIndex() throws Exception { "hello-java, 42, java-new", "python-only, 7, python-only"); } + @Test + @EnabledIfSystemProperty(named = "run.e2e.tests", matches = "true") + public void testReadPythonCompositeDatePartition() throws Exception { + for (boolean legacyName : Arrays.asList(true, false)) { + String suffix = legacyName ? "legacy" : "canonical"; + FileStoreTable table = + (FileStoreTable) + catalog.getTable(identifier("composite_date_java_to_python_" + suffix)); + List rows = + getResult( + table.newRead(), + table.newScan().plan().splits(), + row -> + row.getInt(0) + + ":" + + row.getInt(1) + + ":" + + row.getString(2).toString()); + assertThat(rows).containsExactlyInAnyOrder("1:1:a/b", "2:1:a/b"); + + table.rollbackTo(1L); + List rowsAfterRollback = + getResult( + table.newRead(), + table.newScan().plan().splits(), + row -> + row.getInt(0) + + ":" + + row.getInt(1) + + ":" + + row.getString(2).toString()); + assertThat(rowsAfterRollback).containsExactly("1:1:a/b"); + } + } + private int assignDynamicBucket(FileStoreTable table, InternalRow row) { InternalRowSerializer serializer = new InternalRowSerializer(DataTypes.STRING(), DataTypes.BIGINT()); diff --git a/paimon-python/dev/run_mixed_tests.sh b/paimon-python/dev/run_mixed_tests.sh index 84a54e21dbae..9b8784b973c1 100755 --- a/paimon-python/dev/run_mixed_tests.sh +++ b/paimon-python/dev/run_mixed_tests.sh @@ -99,6 +99,7 @@ run_batched_java_write_tests() { local core_tests="org.apache.paimon.JavaPyE2ETest#testJavaWriteReadPkTable" core_tests="${core_tests}+testJavaWriteDynamicBucketHashIndex" + core_tests="${core_tests}+testJavaWriteCompositeDatePartition" core_tests="${core_tests}+testPKDeletionVectorWrite" core_tests="${core_tests}+testBtreeIndexWrite" core_tests="${core_tests}+testBtreeRawFallbackWrite" @@ -187,7 +188,7 @@ run_java_write_test() { echo "Running Maven test for JavaPyE2ETest.testJavaWriteReadPkTable (Parquet/Orc/Avro)..." echo "Note: Maven may download dependencies on first run, this may take a while..." local parquet_result=0 - if mvn test -Dtest=org.apache.paimon.JavaPyE2ETest#testJavaWriteReadPkTable+testJavaWriteDynamicBucketHashIndex -pl paimon-core -Drun.e2e.tests=true; then + if mvn test -Dtest=org.apache.paimon.JavaPyE2ETest#testJavaWriteReadPkTable+testJavaWriteDynamicBucketHashIndex+testJavaWriteCompositeDatePartition -pl paimon-core -Drun.e2e.tests=true; then echo -e "${GREEN}✓ Java write Parquet/Orc/Avro test completed successfully${NC}" else echo -e "${RED}✗ Java write Parquet/Orc/Avro test failed${NC}" @@ -225,7 +226,7 @@ run_python_read_test() { # Run the parameterized Python test method (runs for both Parquet/Orc/Avro and Lance) echo "Running Python test for JavaPyReadWriteTest.test_read_pk_table..." - if python -m pytest java_py_read_write_test.py::JavaPyReadWriteTest -k "test_read_pk_table or test_read_java_dynamic_bucket_hash_index" -v; then + if python -m pytest java_py_read_write_test.py::JavaPyReadWriteTest -k "test_read_pk_table or test_read_java_dynamic_bucket_hash_index or test_read_java_composite_date_partition" -v; then echo -e "${GREEN}✓ Python test completed successfully${NC}" # source deactivate return 0 @@ -244,7 +245,7 @@ run_python_write_test() { # Run the parameterized Python test method for writing data (pk table, includes bucket num assertion) echo "Running Python test for JavaPyReadWriteTest (test_py_write_read_pk_table)..." - if python -m pytest java_py_read_write_test.py::JavaPyReadWriteTest -k "test_py_write_read_pk_table or test_py_write_dynamic_bucket_hash_index or test_py_write_floating_sequence" -v; then + if python -m pytest java_py_read_write_test.py::JavaPyReadWriteTest -k "test_py_write_read_pk_table or test_py_write_dynamic_bucket_hash_index or test_py_write_floating_sequence or test_py_write_composite_date_partition" -v; then echo -e "${GREEN}✓ Python write test completed successfully${NC}" return 0 else @@ -263,7 +264,7 @@ run_java_read_test() { echo "Running Maven test for JavaPyE2ETest.testReadPkTable (Java Read Parquet/Orc/Avro)..." echo "Note: Maven may download dependencies on first run, this may take a while..." local parquet_result=0 - if mvn test -Dtest=org.apache.paimon.JavaPyE2ETest#testReadPkTable+testReadPythonDynamicBucketHashIndex+testReadPythonFloatingSequence -pl paimon-core -Drun.e2e.tests=true -Dpython.version="$PYTHON_VERSION"; then + if mvn test -Dtest=org.apache.paimon.JavaPyE2ETest#testReadPkTable+testReadPythonDynamicBucketHashIndex+testReadPythonFloatingSequence+testReadPythonCompositeDatePartition -pl paimon-core -Drun.e2e.tests=true -Dpython.version="$PYTHON_VERSION"; then echo -e "${GREEN}✓ Java read Parquet/Orc/Avro test completed successfully${NC}" else echo -e "${RED}✗ Java read Parquet/Orc/Avro test failed${NC}" diff --git a/paimon-python/pypaimon/read/scanner/data_evolution_split_generator.py b/paimon-python/pypaimon/read/scanner/data_evolution_split_generator.py index be995c5acdcd..1efa3fc93105 100644 --- a/paimon-python/pypaimon/read/scanner/data_evolution_split_generator.py +++ b/paimon-python/pypaimon/read/scanner/data_evolution_split_generator.py @@ -128,15 +128,23 @@ def _build_split_from_pack_for_data_evolution( raw_convertible is True only when each range (pack) contains exactly one file. """ splits = [] + data_files = [ + data_file + for file_group in flatten_packed_files + for data_file in file_group + ] + self._set_data_file_paths( + data_files, + file_entries[0].partition, + file_entries[0].bucket, + ) + for i, file_group in enumerate(flatten_packed_files): # In Java: rawConvertible = f.stream().allMatch(file -> file.size() == 1) # This means raw_convertible is True only when each range contains exactly one file pack = packed_files[i] if i < len(packed_files) else [] raw_convertible = all(len(sub_pack) == 1 for sub_pack in pack) - self._set_data_file_paths( - file_group, file_entries[0].partition, file_entries[0].bucket) - if file_group: # Get deletion files for this split data_deletion_files = None diff --git a/paimon-python/pypaimon/read/scanner/split_generator.py b/paimon-python/pypaimon/read/scanner/split_generator.py index 2aa8f2ffb918..d0c4ce14421d 100644 --- a/paimon-python/pypaimon/read/scanner/split_generator.py +++ b/paimon-python/pypaimon/read/scanner/split_generator.py @@ -97,6 +97,21 @@ def _build_split_from_pack( splits = [] if not packed_files or not file_entries: return splits + + # All packed groups are derived from the same partition and bucket. + # Resolve every file together so path compatibility needs at most one + # canonical lookup and one historical fallback per partition/bucket. + data_files = [ + data_file + for file_group in packed_files + for data_file in file_group + ] + self._set_data_file_paths( + data_files, + file_entries[0].partition, + file_entries[0].bucket, + ) + for file_group in packed_files: if use_optimized_path: raw_convertible = True @@ -105,9 +120,6 @@ def _build_split_from_pack( else: raw_convertible = True - self._set_data_file_paths( - file_group, file_entries[0].partition, file_entries[0].bucket) - if file_group: # Get deletion files for this split data_deletion_files = None @@ -130,15 +142,87 @@ def _build_split_from_pack( return splits def _set_data_file_paths(self, files, partition: GenericRow, bucket: int): - """Resolve Java's partition paths without probing the primary file. + """解析规范分区路径,并兼容旧 Python 分区目录。 - A row-sidecar read may not open the primary file at all. Its aligned - path must still be correct when that primary file is unavailable. + 仅当旧、新分区目录不同时才检查文件是否存在。读取 row sidecar 时可能 + 不会打开主数据文件,因此主文件和 sidecar 都要参与目录判定。 """ values = tuple(partition.values) + path_factory = self.table.path_factory() + canonical_bucket = path_factory.bucket_path(values, bucket, canonical_partition=True) + legacy_bucket = path_factory.bucket_path(values, bucket) + + if canonical_bucket == legacy_bucket: + for data_file in files: + if data_file.external_path: + data_file.file_path = data_file.physical_path() + else: + data_file.file_path = canonical_data_file_path( + self.table, values, bucket, data_file.file_name) + return + + path_candidates = [] + for data_file in files: - data_file.file_path = data_file.physical_path() if data_file.external_path else canonical_data_file_path( + if data_file.external_path: + data_file.file_path = data_file.physical_path() + continue + + canonical_path = canonical_data_file_path( self.table, values, bucket, data_file.file_name) + data_file.file_path = canonical_path + canonical_paths = [data_file.physical_path()] + canonical_paths.extend( + data_file.aligned_file_path(name, canonical_bucket) + for name in data_file.extra_files + ) + + data_file.set_file_path( + self.table.table_path, + partition, + bucket, + self.default_part_value, + self.table.options.data_file_path_directory(), + ) + legacy_path = data_file.file_path + legacy_paths = [data_file.physical_path()] + legacy_paths.extend( + data_file.aligned_file_path(name, legacy_bucket) + for name in data_file.extra_files + ) + path_candidates.append(( + data_file, + canonical_path, + legacy_path, + canonical_paths, + legacy_paths, + )) + data_file.file_path = canonical_path + + if not path_candidates: + return + + canonical_paths = { + path for _, _, _, paths, _ in path_candidates for path in paths + } + existing_canonical_paths = self.table.file_io.exists_batch(list(canonical_paths)) + legacy_candidates = [] + + for candidates in path_candidates: + data_file, _, _, canonical_paths, legacy_paths = candidates + if any(existing_canonical_paths.get(path, False) for path in canonical_paths): + continue + legacy_candidates.extend(legacy_paths) + + if not legacy_candidates: + return + + existing_legacy_paths = self.table.file_io.exists_batch(list(set(legacy_candidates))) + for data_file, _, legacy_path, canonical_paths, legacy_paths in path_candidates: + if any(existing_canonical_paths.get(path, False) for path in canonical_paths): + continue + if any(existing_legacy_paths.get(path, False) for path in legacy_paths): + data_file.file_path = legacy_path def _get_deletion_files_for_split( self, diff --git a/paimon-python/pypaimon/tests/date_partition_path_test.py b/paimon-python/pypaimon/tests/date_partition_path_test.py new file mode 100644 index 000000000000..0bf91ff7e5ac --- /dev/null +++ b/paimon-python/pypaimon/tests/date_partition_path_test.py @@ -0,0 +1,360 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Verify DATE partition paths match Java and preserve historical reads.""" + +from datetime import date +from pathlib import Path +from unittest.mock import patch + +import pyarrow as pa +import pytest + +from pypaimon.common.identifier import Identifier +from pypaimon.common.options.core_options import CoreOptions +from pypaimon.filesystem.local_file_io import LocalFileIO +from pypaimon.schema.data_types import AtomicType +from pypaimon.schema.schema import Schema +from pypaimon.schema.schema_manager import SchemaManager +from pypaimon.table.file_store_table import FileStoreTable +from pypaimon.utils.file_store_path_factory import FileStorePathFactory +from pypaimon.write.writer.data_writer import DataWriter + + +def _create_table( + tmp_path, legacy_partition_name=True, composite_partition=False, table_options=None +): + table_path = str(tmp_path / "table") + file_io = LocalFileIO() + fields = [("id", pa.int64()), ("day", pa.date32())] + partition_keys = ["day"] + if composite_partition: + fields.append(("region", pa.string())) + partition_keys.append("region") + arrow_schema = pa.schema(fields) + options = { + "partition.legacy-name": str(legacy_partition_name).lower(), + "scan.native-plan.enabled": "false", + "write.native.enabled": "false", + "commit.native.enabled": "false", + } + options.update(table_options or {}) + schema = Schema.from_pyarrow_schema( + arrow_schema, + partition_keys=partition_keys, + options=options, + ) + table_schema = SchemaManager(file_io, table_path).create_table(schema) + table = FileStoreTable( + file_io, + Identifier.create("default", "t"), + table_path, + table_schema, + ) + return table, arrow_schema + + +def _write_rows(table, arrow_schema, rows): + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(pa.table( + { + name: [row[name] for row in rows] + for name in arrow_schema.names + }, + schema=arrow_schema, + )) + commit.commit(writer.prepare_commit()) + finally: + writer.close() + commit.close() + + +def _write_one_row(table, arrow_schema): + _write_rows(table, arrow_schema, [{"id": 1, "day": date(1970, 1, 2)}]) + + +def _write_legacy_composite_row(table, arrow_schema, row): + historical_bucket = ( + Path(table.table_path) + / "day=1970-01-02" + / "region=a" + / "b" + / "bucket-0" + ) + with patch.object( + DataWriter, + "_generate_file_path", + lambda writer, file_name: str(historical_bucket / file_name), + ): + _write_rows(table, arrow_schema, [row]) + + +def _write_legacy_date_row(table, arrow_schema, row): + historical_bucket = Path(table.table_path) / "day=1970-01-02" / "bucket-0" + with patch.object( + DataWriter, + "_generate_file_path", + lambda writer, file_name: str(historical_bucket / file_name), + ): + _write_rows(table, arrow_schema, [row]) + + +def _read_rows(table): + builder = table.new_read_builder() + splits = builder.new_scan().plan().splits() + result = builder.new_read().to_arrow(splits) + columns = result.to_pydict() + return [ + {name: columns[name][index] for name in result.column_names} + for index in range(result.num_rows) + ] + + +def test_legacy_date_partition_write_uses_java_epoch_day(tmp_path): + table, arrow_schema = _create_table(tmp_path) + + _write_one_row(table, arrow_schema) + + partition_directories = sorted( + path.name for path in Path(table.table_path).glob("day=*") + ) + assert partition_directories == ["day=1"] + assert _read_rows(table) == [{"id": 1, "day": date(1970, 1, 2)}] + + +@pytest.mark.python_plan +@pytest.mark.python_read +@pytest.mark.python_write +@pytest.mark.python_commit +def test_legacy_python_date_partition_directory_remains_readable(tmp_path): + table, arrow_schema = _create_table(tmp_path) + + _write_legacy_date_row( + table, arrow_schema, {"id": 1, "day": date(1970, 1, 2)} + ) + historical_directory = Path(table.table_path) / "day=1970-01-02" + assert historical_directory.joinpath("bucket-0").is_dir() + + assert _read_rows(table) == [{"id": 1, "day": date(1970, 1, 2)}] + + +def test_abort_removes_serialized_date_partition_file(tmp_path): + from pypaimon.write.commit_message_serializer import ( + deserialize_commit_message, + serialize_commit_message, + ) + + table, arrow_schema = _create_table(tmp_path) + builder = table.new_batch_write_builder() + writer = builder.new_write() + try: + writer.write_arrow(pa.Table.from_pylist( + [{"id": 1, "day": date(1970, 1, 2)}], schema=arrow_schema + )) + messages = writer.prepare_commit() + finally: + writer.close() + + file_paths = [ + file.file_path + for message in messages + for file in message.new_files + ] + assert file_paths + assert all(table.file_io.exists(path) for path in file_paths) + + restored_messages = [ + deserialize_commit_message( + serialize_commit_message(message, table.partition_keys_fields), + table.partition_keys_fields, + ) + for message in messages + ] + assert all( + file.file_path is None + for message in restored_messages + for file in message.new_files + ) + + commit = builder.new_commit() + try: + commit.abort(restored_messages) + finally: + commit.close() + + assert not any(table.file_io.exists(path) for path in file_paths) + + +def test_non_legacy_date_partition_write_keeps_iso_name(tmp_path): + table, arrow_schema = _create_table(tmp_path, legacy_partition_name=False) + + _write_one_row(table, arrow_schema) + + partition_directories = sorted( + path.name for path in Path(table.table_path).glob("day=*") + ) + assert partition_directories == ["day=1970-01-02"] + + +def test_legacy_composite_date_write_reads_back_from_canonical_path(tmp_path): + table, arrow_schema = _create_table(tmp_path, composite_partition=True) + + _write_rows(table, arrow_schema, [{ + "id": 1, "day": date(1970, 1, 2), "region": "a/b", + }]) + + assert _read_rows(table) == [{ + "id": 1, "day": date(1970, 1, 2), "region": "a/b", + }] + assert (Path(table.table_path) / "day=1" / "region=a%2Fb" / "bucket-0").is_dir() + + +@pytest.mark.python_plan +@pytest.mark.python_read +@pytest.mark.python_write +@pytest.mark.python_commit +def test_old_and_new_composite_partition_files_coexist_after_append(tmp_path): + table, arrow_schema = _create_table(tmp_path, composite_partition=True) + row = {"day": date(1970, 1, 2), "region": "a/b"} + + _write_legacy_composite_row(table, arrow_schema, {"id": 1, **row}) + first_snapshot_id = table.snapshot_manager().get_latest_snapshot().id + historical_bucket = ( + Path(table.table_path) + / "day=1970-01-02" + / "region=a" + / "b" + / "bucket-0" + ) + assert historical_bucket.is_dir() + + _write_rows(table, arrow_schema, [{"id": 2, **row}]) + canonical_directory = Path(table.table_path) / "day=1" / "region=a%2Fb" + + assert {item["id"] for item in _read_rows(table)} == {1, 2} + assert canonical_directory.joinpath("bucket-0").is_dir() + + table.rollback_to(first_snapshot_id) + + assert [item["id"] for item in _read_rows(table)] == [1] + + +@pytest.mark.python_plan +@pytest.mark.python_read +@pytest.mark.python_write +@pytest.mark.python_commit +@pytest.mark.parametrize("historical_layout", [False, True]) +def test_composite_date_path_lookups_are_batched_across_split_packs( + tmp_path, monkeypatch, historical_layout +): + table, arrow_schema = _create_table( + tmp_path, + composite_partition=True, + table_options={ + CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key(): "1b", + CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key(): "1b", + }, + ) + row = {"day": date(1970, 1, 2), "region": "a/b"} + + for row_id in (1, 2): + data = {"id": row_id, **row} + if historical_layout: + _write_legacy_composite_row(table, arrow_schema, data) + else: + _write_rows(table, arrow_schema, [data]) + + lookup_batches = [] + file_io = table.file_io + exists_batch = file_io.exists_batch + + def track_exists_batch(paths): + lookup_batches.append(list(paths)) + return exists_batch(paths) + + monkeypatch.setattr(file_io, "exists_batch", track_exists_batch) + splits = table.new_read_builder().new_scan().plan().splits() + + assert len(splits) == 2 + assert sum(len(split.files) for split in splits) == 2 + expected_batch_sizes = [2, 2] if historical_layout else [2] + assert [len(paths) for paths in lookup_batches] == expected_batch_sizes + assert {item["id"] for item in _read_rows(table)} == {1, 2} + + +def test_non_legacy_composite_date_write_uses_canonical_escaping(tmp_path): + table, arrow_schema = _create_table( + tmp_path, legacy_partition_name=False, composite_partition=True) + + _write_rows(table, arrow_schema, [{ + "id": 1, "day": date(1970, 1, 2), "region": "a/b", + }]) + + assert _read_rows(table) == [{ + "id": 1, "day": date(1970, 1, 2), "region": "a/b", + }] + assert (Path(table.table_path) / "day=1970-01-02" / "region=a%2Fb" / "bucket-0").is_dir() + + +def test_date_compatibility_does_not_reformat_other_partition_types(): + factory = FileStorePathFactory( + "/table", + ["day", "region"], + "__DEFAULT_PARTITION__", + "parquet", + "data-", + "changelog-", + True, + False, + None, + partition_types=[AtomicType("DATE"), AtomicType("STRING")], + ) + + assert factory.relative_bucket_path( + (date(1970, 1, 2), "a/b"), 0, canonical_partition=True + ) == "day=1/region=a%2Fb/bucket-0" + + +def test_date_compatibility_applies_to_external_data_and_bucket_index_paths(): + partition = (date(1970, 1, 2),) + factory = FileStorePathFactory( + "/table", + ["day"], + "__DEFAULT_PARTITION__", + "parquet", + "data-", + "changelog-", + True, + False, + None, + external_paths=["/external"], + external_path_strategy="round-robin", + index_file_in_data_file_dir=True, + partition_types=[AtomicType("DATE")], + ) + + external_provider = factory.create_external_path_provider( + partition, 0, canonical_partition=True) + assert external_provider.get_next_external_data_path("data.parquet") == ( + "/external/day=1/bucket-0/data.parquet" + ) + assert factory.new_bucket_index_path(partition, 0, "index-1") == ( + "/external/day=1/bucket-0/index-1", + True, + ) diff --git a/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py b/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py index fa6ff03103bc..afa6e26ab7db 100644 --- a/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py +++ b/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py @@ -142,6 +142,49 @@ def test_py_write_dynamic_bucket_hash_index(self): writer.close() commit.close() + def test_read_java_composite_date_partition(self): + for legacy_name in (True, False): + suffix = 'legacy' if legacy_name else 'canonical' + table = self.catalog.get_table( + 'default.composite_date_java_to_python_' + suffix) + read_builder = table.new_read_builder() + result = read_builder.new_read().to_arrow( + read_builder.new_scan().plan().splits()) + self.assertEqual(result.to_pydict(), { + 'id': [1], + 'day': [datetime.date(1970, 1, 2)], + 'region': ['a/b'], + }) + + def test_py_write_composite_date_partition(self): + arrow_schema = pa.schema([ + ('id', pa.int32()), + ('day', pa.date32()), + ('region', pa.string()), + ]) + for legacy_name in (True, False): + suffix = 'legacy' if legacy_name else 'canonical' + table = self.catalog.get_table( + 'default.composite_date_java_to_python_' + suffix) + builder = table.new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(pa.table({ + 'id': [2], + 'day': [datetime.date(1970, 1, 2)], + 'region': ['a/b'], + }, schema=arrow_schema)) + commit.commit(writer.prepare_commit()) + finally: + writer.close() + commit.close() + + read_builder = table.new_read_builder() + result = read_builder.new_read().to_arrow( + read_builder.new_scan().plan().splits()) + self.assertEqual( + set(result.to_pydict()['id']), {1, 2}) + @parameterized.expand([ (type_name, order, grouping) for type_name in ('float', 'double')