This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new b026557f6d Expand Python native plan and read coverage for BLOB and PK
tables (#10065)
b026557f6d is described below
commit b026557f6d8aa4bab27bfeb1ea1e0a00224fa845
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 21:31:42 2026 +0800
Expand Python native plan and read coverage for BLOB and PK tables (#10065)
---
paimon-python/pypaimon/read/table_read.py | 37 +++-----
.../pypaimon/tests/native_plan_integration_test.py | 103 ++++++++++++++++++++-
.../tests/native_plan_materialized_pk_test.py | 83 +++++++++++++++++
.../pypaimon/tests/native_plan_rest_test.py | 47 ++++++++++
paimon-python/pypaimon/tests/native_read_test.py | 53 +++++++----
5 files changed, 283 insertions(+), 40 deletions(-)
diff --git a/paimon-python/pypaimon/read/table_read.py
b/paimon-python/pypaimon/read/table_read.py
index aa6f395229..1b3a4585d3 100644
--- a/paimon-python/pypaimon/read/table_read.py
+++ b/paimon-python/pypaimon/read/table_read.py
@@ -420,8 +420,10 @@ class TableRead:
return None
if not splits:
return []
+ if not self._native_blob_view_supported():
+ return None
if (self._deferred_blob_limit_may_prune(splits)
- and not self._native_pruning_blob_limit_supported()):
+ and not self.table.options.data_evolution_enabled()):
return None
# Query authorization has additional filtering, masking and projection
# semantics which are already implemented by the Python reader.
@@ -922,28 +924,19 @@ class TableRead:
or self._native_inline_blob_fields())
and not self._limit_covers_all_splits(splits))
- def _native_pruning_blob_limit_supported(self) -> bool:
- """Rust caps DE batches before payload resolution when no post-filter
is needed.
-
- A predicate on a managed BLOB or BLOB view can require payload I/O
- before the output quota is known. Inline descriptors can still use
- native reads if the predicate only references ordinary columns.
- """
- if not self.table.options.data_evolution_enabled():
- return False
+ def _native_blob_view_supported(self) -> bool:
+ """Use native view resolution only with the REST catalog
environment."""
read_names = {field.name for field in self._scan_read_type}
- if self.table.options.blob_view_fields() & read_names:
- return False
- if self.predicate is not None:
- # Managed BLOBs are decoded by the physical file reader, before a
- # residual filter. Inline descriptors are resolved later, so a
- # predicate on ordinary columns can safely select rows first.
- if self._deferred_blob_fields:
- return False
- from pypaimon.read.push_down_utils import predicate_field_names
- if predicate_field_names(self.predicate) &
self._native_inline_blob_fields():
- return False
- return True
+ view_fields = self.table.options.blob_view_fields() & read_names
+ if not view_fields or not
self.table.options.blob_view_resolve_enabled():
+ return True
+ loader = getattr(getattr(self.table, 'catalog_environment', None),
+ 'catalog_loader', None)
+ if loader is None:
+ # Python also leaves view structs unresolved without a loader.
+ return True
+ from pypaimon.read.native_plan import _catalog_metastore
+ return _catalog_metastore(loader) == 'rest'
def _native_inline_blob_fields(self) -> set:
"""Return configured BLOB fields that native reads resolve eagerly."""
diff --git a/paimon-python/pypaimon/tests/native_plan_integration_test.py
b/paimon-python/pypaimon/tests/native_plan_integration_test.py
index a053ed5627..a3caec88f8 100644
--- a/paimon-python/pypaimon/tests/native_plan_integration_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_integration_test.py
@@ -37,7 +37,7 @@ from pypaimon.read.native_plan import (
native_split_from_python,
)
from pypaimon.schema.schema_change import SchemaChange
-from pypaimon.table.row.blob import BlobDescriptor
+from pypaimon.table.row.blob import BlobDescriptor, BlobViewStruct
from pypaimon.utils.range import Range
@@ -126,6 +126,35 @@ class NativePlanIntegrationTest(unittest.TestCase):
side_effect=AssertionError('Python reader was used')):
return builder.new_read().to_arrow(plan.splits()).to_pylist()
+ def test_row_id_predicate_uses_native_plan_and_read(self):
+ for data_evolution in (False, True):
+ with self.subTest(data_evolution=data_evolution):
+ name = 'row_id_predicate_de' if data_evolution else
'row_id_predicate_append'
+ self.cat.create_table('default.' + name,
Schema.from_pyarrow_schema(
+ self.schema, options={
+ 'row-tracking.enabled': 'true',
+ 'data-evolution.enabled': str(data_evolution).lower(),
+ }), False)
+ self._write(name, [
+ {'k': 1, 'v': 'a'}, {'k': 2, 'v': 'b'}, {'k': 3, 'v': 'c'},
+ ])
+ table = self.cat.get_table('default.' + name).copy({
+ 'scan.native-plan.enabled': 'true', 'read.native.enabled':
'true',
+ })
+ predicates = (table.new_read_builder()
+ .with_projection(['k', '_ROW_ID'])
+ .new_predicate_builder())
+ builder = table.new_read_builder().with_projection(['k',
'_ROW_ID'])
+ builder.with_filter(predicates.and_predicates([
+ predicates.equal('_ROW_ID', 1), predicates.equal('k', 2),
+ ]))
+ plan = builder.new_scan().plan()
+ self.assertTrue(builder.explain().native_planned)
+ self.assertTrue(all(getattr(split, '_native_split', None) is
not None
+ for split in plan.splits()))
+ self.assertEqual(self._native_rows(builder, plan),
+ [{'k': 2, '_ROW_ID': 1}])
+
def test_primary_key_matches_normal_plan(self):
self.cat.create_table('default.pk_t', Schema.from_pyarrow_schema(
self.schema, primary_keys=['k'], options={'bucket': '1'}), False)
@@ -1027,6 +1056,78 @@ class NativePlanIntegrationTest(unittest.TestCase):
native.assert_called_once()
self.assertEqual(result.to_pydict(), {'id': [1], 'payload': [b'a']})
+ @unittest.skipUnless(native_reader_available(),
+ "pypaimon-rust native reader API not installed")
+ def test_native_read_pruning_blob_limit_with_predicates(self):
+ schema = pa.schema([('id', pa.int32()), ('payload',
pa.large_binary())])
+ self.cat.create_table('default.native_blob_predicate_t',
+ Schema.from_pyarrow_schema(schema, options={
+ 'row-tracking.enabled': 'true',
+ 'data-evolution.enabled': 'true',
+ }), False)
+ table = self.cat.get_table('default.native_blob_predicate_t')
+ writer = table.new_batch_write_builder().new_write()
+ writer.write_arrow(pa.Table.from_pydict({
+ 'id': [1, 2, 3], 'payload': [b'skip', b'selected', b'late'],
+ }, schema=schema))
+
table.new_batch_write_builder().new_commit().commit(writer.prepare_commit())
+ writer.close()
+
+ table = table.copy({
+ 'scan.native-plan.enabled': 'true', 'read.native.enabled': 'true',
+ })
+ for field, value in (('id', 2), ('payload', b'selected')):
+ with self.subTest(predicate_field=field):
+ builder = table.new_read_builder().with_limit(1)
+
builder.with_filter(builder.new_predicate_builder().equal(field, value))
+ plan = builder.new_scan().plan()
+ self.assertTrue(builder.explain().native_planned)
+ self.assertTrue(all(getattr(split, '_native_split', None) is
not None
+ for split in plan.splits()))
+ self.assertEqual(self._native_rows(builder, plan),
+ [{'id': 2, 'payload': b'selected'}])
+
+ @unittest.skipUnless(native_reader_available(),
+ "pypaimon-rust native reader API not installed")
+ def test_filesystem_blob_view_without_limit_stays_on_python_reader(self):
+ schema = pa.schema([('id', pa.int32()), ('payload',
pa.large_binary())])
+ options = {'row-tracking.enabled': 'true', 'data-evolution.enabled':
'true'}
+ self.cat.create_table('default.fs_blob_source',
+ Schema.from_pyarrow_schema(schema,
options=options), False)
+ source = self.cat.get_table('default.fs_blob_source')
+ writer = source.new_batch_write_builder().new_write()
+ writer.write_arrow(pa.Table.from_pydict({
+ 'id': [1], 'payload': [b'from-source'],
+ }, schema=schema))
+
source.new_batch_write_builder().new_commit().commit(writer.prepare_commit())
+ writer.close()
+ payload_id = next(field.id for field in source.table_schema.fields
+ if field.name == 'payload')
+
+ self.cat.create_table('default.fs_blob_view',
Schema.from_pyarrow_schema(
+ schema, options=dict(options, **{'blob-view-field': 'payload'})),
False)
+ view = self.cat.get_table('default.fs_blob_view')
+ writer = view.new_batch_write_builder().new_write()
+ writer.write_arrow(pa.Table.from_pydict({
+ 'id': [10],
+ 'payload': [BlobViewStruct(
+ 'default.fs_blob_source', payload_id, 0).serialize()],
+ }, schema=schema))
+
view.new_batch_write_builder().new_commit().commit(writer.prepare_commit())
+ writer.close()
+
+ view = view.copy({
+ 'scan.native-plan.enabled': 'true', 'read.native.enabled': 'true',
+ })
+ builder = view.new_read_builder()
+ plan = builder.new_scan().plan()
+ self.assertTrue(builder.explain().native_planned)
+ with patch('pypaimon.read.native_plan.native_read',
+ side_effect=AssertionError('native view read should not
run')) as native:
+
self.assertEqual(builder.new_read().to_arrow(plan.splits()).to_pylist(),
+ [{'id': 10, 'payload': b'from-source'}])
+ native.assert_not_called()
+
@unittest.skipUnless(native_reader_available(),
"pypaimon-rust native reader API not installed")
def test_native_read_pruning_limit_defers_descriptor_blob_payload_io(self):
diff --git a/paimon-python/pypaimon/tests/native_plan_materialized_pk_test.py
b/paimon-python/pypaimon/tests/native_plan_materialized_pk_test.py
index c418e71173..d7595757be 100644
--- a/paimon-python/pypaimon/tests/native_plan_materialized_pk_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_materialized_pk_test.py
@@ -213,3 +213,86 @@ def
test_first_row_level_zero_merges_before_filtering(tmp_path, target_size, com
actual = builder.new_read().to_arrow(
plan.splits()).column('id').to_pylist()
assert sorted(actual) == expected
+
+
[email protected]('engine,merged_value', [
+ ('partial-update', 5), ('aggregation', 15),
+])
+def test_pk_merge_engines_read_deletion_vectors_natively(tmp_path, engine,
merged_value):
+ schema = pa.schema([('id', pa.int64()), ('value', pa.int64())])
+ options = {
+ 'bucket': '1', 'merge-engine': engine, 'file.format': 'parquet',
+ 'deletion-vectors.enabled': 'true', 'deletion-vectors.merge-on-read':
'true',
+ }
+ if engine == 'aggregation':
+ options['fields.value.aggregate-function'] = 'sum'
+ catalog = CatalogFactory.create({'warehouse': str(tmp_path)})
+ catalog.create_database('default', True)
+ catalog.create_table('default.t', Schema.from_pyarrow_schema(
+ schema, primary_keys=['id'], options=options), False)
+ table = catalog.get_table('default.t')
+ # Each run has unique keys, so the Python writer's aggregation fallback
+ # leaves its physical rows intact; aggregation happens across the runs.
+ first_file = None
+ for rows, level in (([{'id': 1, 'value': 10}, {'id': 2, 'value': 20}], 1),
+ ([{'id': 1, 'value': 5}, {'id': 3, 'value': 30}], 0)):
+ builder = table.new_batch_write_builder()
+ writer, commit = builder.new_write(), builder.new_commit()
+ try:
+ writer.write_arrow(pa.Table.from_pylist(rows, schema=schema))
+ messages = writer.prepare_commit()
+ for message in messages:
+ message.new_files = [replace(file, level=level)
+ for file in message.new_files]
+ if first_file is None:
+ first_file = message.new_files[0]
+ commit.commit(messages)
+ finally:
+ writer.close()
+ commit.close()
+
+ first_path = table.path_factory().bucket_path((), 0) + '/' +
first_file.file_name
+ deleted_position =
pq.read_table(first_path).column('id').to_pylist().index(2)
+ vector = BitmapDeletionVector()
+ vector.delete(deleted_position)
+ entry = TableDeleteByRowId(table)._write_deletion_vector_index(
+ GenericRow([], []), 0, {first_file.file_name: vector})
+ commit = table.new_batch_write_builder().new_commit()
+ try:
+ commit.commit([CommitMessage(partition=(), bucket=0, new_files=[],
index_adds=[entry])])
+ finally:
+ commit.close()
+
+ expected = [{'id': 1, 'value': merged_value}, {'id': 3, 'value': 30}]
+ for predicate_value, projection in ((None, None), (merged_value, ['id'])):
+ for native in (False, True):
+ candidate = table.copy({
+ 'scan.native-plan.enabled': str(native).lower(),
+ 'read.native.enabled': str(native).lower(),
+ })
+ predicates = candidate.new_read_builder().new_predicate_builder()
+ builder = candidate.new_read_builder()
+ if projection is not None:
+ builder.with_projection(projection)
+ if predicate_value is not None:
+ builder.with_filter(predicates.equal('value', predicate_value))
+ scan = builder.new_scan()
+ if native:
+ with patch.object(scan.file_scanner, 'scan',
+ side_effect=AssertionError('native plan fell
back')):
+ plan = scan.plan()
+ assert all(getattr(split, '_native_split', None) is not None
+ for split in plan.splits())
+ guard =
patch('pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError('native read fell
back'))
+ else:
+ plan = scan.plan()
+ guard = ExitStack()
+ with guard:
+ actual = builder.new_read().to_arrow(plan.splits()).to_pylist()
+ wanted = (expected if predicate_value is None
+ else [{'id': 1}])
+ assert sorted(actual, key=lambda row: row['id']) == wanted
+ assert any(deletion is not None
+ for split in plan.splits()
+ for deletion in split.data_deletion_files or [])
diff --git a/paimon-python/pypaimon/tests/native_plan_rest_test.py
b/paimon-python/pypaimon/tests/native_plan_rest_test.py
index 2a34ce4a29..a2c2a92fde 100644
--- a/paimon-python/pypaimon/tests/native_plan_rest_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_rest_test.py
@@ -25,6 +25,7 @@ from pypaimon.api.api_response import ConfigResponse,
ErrorResponse, GetTableSna
from pypaimon.api.auth import BearTokenAuthProvider
from pypaimon.read.native_plan import native_runtime_available
from pypaimon.snapshot.table_snapshot import TableSnapshot
+from pypaimon.table.row.blob import BlobViewStruct
from pypaimon.tests.rest.rest_server import RESTCatalogServer
@@ -149,3 +150,49 @@ def
test_rest_dotted_database_and_table_keep_identity(rest_catalog, branch):
_assert_parity(table, [{'id': 1, 'value': 'old'}], 1)
assert all((call.args[1].get_database_name(),
call.args[1].get_table_name())
== ('namespace.database', 'table.with.dots') for call in
load.call_args_list)
+
+
+def
test_rest_blob_view_limit_filters_before_resolving_unselected_view(rest_catalog):
+ catalog, _ = rest_catalog
+ schema = pa.schema([('id', pa.int32()), ('payload', pa.large_binary())])
+ options = {'row-tracking.enabled': 'true', 'data-evolution.enabled':
'true'}
+ catalog.create_table('default.source', Schema.from_pyarrow_schema(
+ schema, options=options), False)
+ source = catalog.get_table('default.source')
+ writer = source.new_batch_write_builder().new_write()
+ writer.write_arrow(pa.Table.from_pydict({
+ 'id': [1, 2], 'payload': [b'first', b'selected'],
+ }, schema=schema))
+
source.new_batch_write_builder().new_commit().commit(writer.prepare_commit())
+ writer.close()
+ payload_id = next(field.id for field in source.table_schema.fields
+ if field.name == 'payload')
+
+ catalog.create_table('default.views', Schema.from_pyarrow_schema(
+ schema, options=dict(options, **{'blob-view-field': 'payload'})),
False)
+ views = catalog.get_table('default.views')
+ writer = views.new_batch_write_builder().new_write()
+ writer.write_arrow(pa.Table.from_pydict({
+ 'id': [10, 11], 'payload': [
+ BlobViewStruct('default.source', payload_id, 99).serialize(),
+ BlobViewStruct('default.source', payload_id, 1).serialize(),
+ ],
+ }, schema=schema))
+
views.new_batch_write_builder().new_commit().commit(writer.prepare_commit())
+ writer.close()
+
+ views = views.copy({
+ 'scan.native-plan.enabled': 'true', 'read.native.enabled': 'true',
+ })
+ builder = views.new_read_builder().with_limit(1)
+ builder.with_filter(builder.new_predicate_builder().equal('id', 11))
+ scan = builder.new_scan()
+ with patch.object(scan.file_scanner, 'scan',
+ side_effect=AssertionError('native view plan fell
back')):
+ plan = scan.plan()
+ assert all(getattr(split, '_native_split', None) is not None
+ for split in plan.splits())
+ with patch('pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError('native view read fell back')):
+ assert builder.new_read().to_arrow(plan.splits()).to_pylist() == [
+ {'id': 11, 'payload': b'selected'}]
diff --git a/paimon-python/pypaimon/tests/native_read_test.py
b/paimon-python/pypaimon/tests/native_read_test.py
index c9649e441f..6a33c32374 100644
--- a/paimon-python/pypaimon/tests/native_read_test.py
+++ b/paimon-python/pypaimon/tests/native_read_test.py
@@ -39,6 +39,7 @@ def _table_read(limit=None):
read.table.options.blob_as_descriptor.return_value = False
read.table.options.blob_descriptor_fields.return_value = set()
read.table.options.blob_view_fields.return_value = set()
+ read.table.options.blob_view_resolve_enabled.return_value = True
read.table.options.data_evolution_enabled.return_value = False
read.predicate = None
read.read_type = [DataField(0, 'id', AtomicType('INT'))]
@@ -825,30 +826,48 @@ def
test_native_read_pruning_blob_limit_uses_data_evolution_reader(descriptor):
native.assert_called_once()
-def test_native_read_pruning_blob_limit_keeps_predicate_and_view_fallback():
[email protected]('layout', ['managed', 'descriptor', 'rest_view'])
+def test_native_read_pruning_blob_limit_allows_predicate(layout):
read = _blob_table_read(limit=1)
read.table.options.data_evolution_enabled.return_value = True
- read._deferred_blob_fields = {'payload'}
- split = _Split('payload.blob')
+ read.predicate = Mock()
+ if layout == 'managed':
+ read._deferred_blob_fields = {'payload'}
+ elif layout == 'descriptor':
+ read.table.options.blob_descriptor_fields.return_value = {'payload'}
+ else:
+ read.table.options.blob_view_fields.return_value = {'payload'}
+ split = _Split('payload.blob' if layout == 'managed' else
'payload.parquet')
split._native_split = object()
split.merged_row_count = Mock(return_value=2)
+ batch = pa.record_batch(
+ [pa.array([b'selected'], type=pa.large_binary())], names=['payload'])
- with patch('pypaimon.read.native_plan.native_read') as native:
- read.predicate = Mock()
- assert read._try_native_batches(
- [split], pa.schema([('payload', pa.large_binary())])) is None
- read.predicate = None
- read.table.options.blob_view_fields.return_value = {'payload'}
+ with patch('pypaimon.read.native_plan._catalog_metastore',
return_value='rest'), \
+ patch('pypaimon.read.native_plan.native_read',
return_value=[batch]) as native:
+ actual = list(read._try_native_batches(
+ [split], pa.schema([('payload', pa.large_binary())])))
+
+ assert actual == [batch]
+ native.assert_called_once()
+ assert native.call_args.kwargs['predicate'] is read.predicate
+ assert native.call_args.kwargs['limit'] == 1
+
+
[email protected]('limit', [None, 1])
+def test_native_read_blob_view_keeps_filesystem_fallback(limit):
+ read = _blob_table_read(limit=limit)
+ read.table.options.data_evolution_enabled.return_value = True
+ read.table.options.blob_view_fields.return_value = {'payload'}
+ read.predicate = Mock()
+ split = _Split('payload.parquet')
+ split._native_split = object()
+ split.merged_row_count = Mock(return_value=2)
+
+ with patch('pypaimon.read.native_plan._catalog_metastore',
return_value='filesystem'), \
+ patch('pypaimon.read.native_plan.native_read') as native:
assert read._try_native_batches(
[split], pa.schema([('payload', pa.large_binary())])) is None
- read.table.options.blob_view_fields.return_value = set()
- read.table.options.blob_descriptor_fields.return_value = {'payload'}
- read._deferred_blob_fields = set()
- read.predicate = Mock()
- with patch('pypaimon.read.push_down_utils.predicate_field_names',
- return_value={'payload'}):
- assert read._try_native_batches(
- [split], pa.schema([('payload', pa.large_binary())])) is None
native.assert_not_called()