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()
 

Reply via email to