Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 71 additions & 1 deletion paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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<String> 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<String> 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());
Expand Down
9 changes: 5 additions & 4 deletions paimon-python/dev/run_mixed_tests.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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}"
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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}"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
98 changes: 91 additions & 7 deletions paimon-python/pypaimon/read/scanner/split_generator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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,
Expand Down
Loading
Loading