Skip to content
Merged
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
11 changes: 8 additions & 3 deletions paimon-python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -161,13 +161,18 @@ finally:
The native writer returns ordinary PyPaimon commit messages, so the Python
committer also works when `commit.native.enabled` is false. Batch overwrite and
reusable stream writers retain the builder's commit user and identifier. Native
write is currently limited to Parquet tables without BLOB fields or
data-evolution mode, on the same filesystem/JDBC publication route as native
commit. Writer methods requiring Python's specialized path select the Python
write supports Parquet append, primary-key and data-evolution tables without
BLOB fields or optional data-evolution row sidecars, on the same filesystem/JDBC
publication route as native commit. Writer methods requiring Python's specialized path select the Python
writer before native data is written. If the runtime or table route is
unavailable, write uses Python. Once Rust starts writing a batch, errors
propagate without retrying that batch through Python.

Native writes honor `data-file.path-directory` and the configured
`data-file.external-paths` strategy. Existing files keep their recorded locations
when the write destinations change. Python and native readers and committers
can exchange these files, including external data files and their index sidecars.

Both native options are disabled by default.

# Native commit
Expand Down
6 changes: 3 additions & 3 deletions paimon-python/pypaimon/common/options/core_options.py
Original file line number Diff line number Diff line change
Expand Up @@ -1692,13 +1692,13 @@ def data_file_external_paths_weights(self, default=None):
)
if value is None:
return None
parts = value.split(",")
parts = value.rstrip(",").split(",")
weights = []
for part in parts:
parsed = int(part.strip())
if parsed <= 0:
if parsed <= 0 or parsed > 2147483647:
raise ValueError(
f"Weight must be positive, got: {parsed}"
f"Weight must be a positive 32-bit integer, got: {parsed}"
)
weights.append(parsed)
return weights
Expand Down
22 changes: 22 additions & 0 deletions paimon-python/pypaimon/manifest/schema/data_file_meta.py
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,28 @@ def create(
file_path=file_path,
)

def physical_path(self):
"""Resolve persisted Java path text for Python file I/O without URL decoding."""
from pypaimon.utils.path import to_file_io_path
path = self.external_path or self.file_path
return to_file_io_path(str(path)) if path is not None else None

def aligned_file_path(self, file_name, bucket_path=None):
"""Place a sidecar beside its data file, including external locations."""
from pypaimon.utils.path import resolve_path, to_file_io_path
data_path = self.physical_path()
parent = data_path.rsplit('/', 1)[0] if data_path and '/' in data_path else bucket_path
return to_file_io_path(resolve_path(parent, file_name) if parent else file_name)

def collect_files(self, bucket_path=None):
"""Return the physical data file and its aligned sidecars for cleanup."""
from pypaimon.utils.path import resolve_path, to_file_io_path
path = self.physical_path()
if not path and bucket_path:
path = to_file_io_path(resolve_path(bucket_path, self.file_name))
return ([path] if path else []) + [
self.aligned_file_path(name, bucket_path) for name in self.extra_files]

def set_file_path(
self, table_path: str, partition: GenericRow, bucket: int,
default_part_value: str = "__DEFAULT_PARTITION__",
Expand Down
11 changes: 3 additions & 8 deletions paimon-python/pypaimon/read/split_read.py
Original file line number Diff line number Diff line change
Expand Up @@ -287,7 +287,7 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool,
read_paimon_predicate = None

# Use external_path if available, otherwise use file_path
file_path = file.external_path if file.external_path else file.file_path
file_path = file.physical_path()
file_format = format_identifier(os.path.basename(file_path))

batch_size = self.table.options.read_batch_size()
Expand Down Expand Up @@ -607,12 +607,7 @@ def _is_row_sidecar_file(file_name: str) -> bool:

@staticmethod
def _aligned_extra_file_path(file: DataFileMeta, extra_file: str) -> str:
if "://" in extra_file or extra_file.startswith("/"):
return extra_file
file_path = file.external_path if file.external_path else file.file_path
if not file_path or "/" not in file_path:
return extra_file
return f"{file_path.rsplit('/', 1)[0]}/{extra_file}"
return file.aligned_file_path(extra_file)

def _get_fields_and_predicate(self, schema_id: int, read_fields):
key = (schema_id, tuple(read_fields))
Expand Down Expand Up @@ -1751,7 +1746,7 @@ def _create_raw_blob_file_reader(
if not row_indices:
return None

file_path = file.external_path if file.external_path else file.file_path
file_path = file.physical_path()
blob_parallelism = self._blob_parallelism
return FormatBlobReader(
self.table.file_io,
Expand Down
2 changes: 1 addition & 1 deletion paimon-python/pypaimon/read/split_serializer.py
Original file line number Diff line number Diff line change
Expand Up @@ -655,7 +655,7 @@ def _datafilemeta_from_row(row_bytes: bytes, bucket_path: str, arity: int,
write_cols_sequences=(
_decode_non_null_long_array(g(20)) if arity >= 21 else None),
)
meta.file_path = external_path if external_path else to_file_io_path(
meta.file_path = meta.physical_path() if external_path else to_file_io_path(
"%s/%s" % (bucket_path.rstrip('/'), file_name))
return meta

Expand Down
8 changes: 8 additions & 0 deletions paimon-python/pypaimon/tests/external_paths_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -265,6 +265,14 @@ def test_mismatched_lengths_raises(self):
class WeightsParsingTest(unittest.TestCase):
"""Test CoreOptions.data_file_external_paths_weights() parsing and validation."""

def test_weights_use_java_integer_range_and_trailing_comma_rules(self):
self.assertEqual(CoreOptions.from_dict({
'data-file.external-paths.weights': '1,2147483647,,'
}).data_file_external_paths_weights(), [1, 2147483647])
for value in ['2147483648', '1,9999999999', '', '1,,2']:
with self.subTest(value=value), self.assertRaises(ValueError):
CoreOptions.from_dict({'data-file.external-paths.weights': value}).data_file_external_paths_weights()

def test_valid_weights(self):
"""Normal comma-separated positive integers."""
from pypaimon.common.options.core_options import CoreOptions
Expand Down
39 changes: 38 additions & 1 deletion paimon-python/pypaimon/tests/file_store_commit_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,13 @@
import uuid
from dataclasses import replace
from datetime import datetime
from pathlib import Path
from tempfile import TemporaryDirectory
from unittest.mock import MagicMock, Mock, patch

from pypaimon.common.options.core_options import CoreOptions
from pypaimon.common.options.options import Options
from pypaimon.filesystem.local_file_io import LocalFileIO
from pypaimon.manifest.schema.data_file_meta import DataFileMeta
from pypaimon.manifest.schema.manifest_entry import ManifestEntry
from pypaimon.manifest.schema.manifest_file_meta import ManifestFileMeta
Expand Down Expand Up @@ -78,15 +81,49 @@ def test_overwrite_validates_message_baseline(self):

class TestAbortCommitMessages(unittest.TestCase):

@staticmethod
def _file_meta(**kwargs):
return DataFileMeta(
file_name='data.parquet', file_size=1, row_count=1,
min_key=None, max_key=None, key_stats=None, value_stats=None,
min_sequence_number=0, max_sequence_number=0, schema_id=0,
level=0, extra_files=kwargs.pop('extra_files', []), **kwargs)

def test_reconstructs_local_path_after_wire_decode(self):
table = Mock()
table.path_factory.return_value.bucket_path.return_value = '/table/p=1/bucket-0'
file = Mock(file_name='data.parquet', external_path=None, file_path=None)
file = self._file_meta()
message = CommitMessage((1,), 0, [file])
_abort_commit_messages(table, [message])
table.file_io.delete_quietly.assert_called_once_with(
'/table/p=1/bucket-0/data.parquet')

def test_reconstructs_aligned_sidecars_after_wire_decode(self):
table = Mock()
table.path_factory.return_value.bucket_path.return_value = '/table/p=1/bucket-0'
file = self._file_meta(extra_files=['data.parquet.index'])
_abort_commit_messages(table, [CommitMessage((1,), 0, [file])])
self.assertEqual(table.file_io.delete_quietly.call_args_list, [
unittest.mock.call('/table/p=1/bucket-0/data.parquet'),
unittest.mock.call('/table/p=1/bucket-0/data.parquet.index'),
])

def test_deletes_literal_external_paths_and_preserves_metadata(self):
with TemporaryDirectory() as directory:
parent = Path(directory) / 'pt=a%2Fb%25%3F%23'
parent.mkdir()
paths = [parent / 'data.parquet', parent / 'data.parquet.index']
for path in paths:
path.touch()
external_path = 'file:' + paths[0].as_posix()
file = self._file_meta(external_path=external_path, extra_files=[paths[1].name])
table = Mock(file_io=LocalFileIO())
_abort_commit_messages(table, [CommitMessage(('a/b%?#',), 0, [file])])
self.assertFalse(any(path.exists() for path in paths))
self.assertEqual(file.external_path, external_path)
self.assertEqual(file.extra_files, [paths[1].name])
table.path_factory.assert_not_called()

def test_index_path_failure_does_not_escape_abort(self):
table = Mock()
table.path_factory.side_effect = RuntimeError("path lookup failed")
Expand Down
Loading
Loading