leaves12138 commented on code in PR #9954:
URL: https://github.com/apache/paimon/pull/9954#discussion_r4045626645


##########
paimon-python/pypaimon/read/table_read.py:
##########
@@ -282,21 +304,199 @@ def to_arrow(
         if self.include_row_kind:
             schema = self._add_row_kind_to_schema(schema)
 
+        native_parallel = self._should_run_parallel(splits, effective)
+        native_bp = effective_bp
+        if native_parallel:
+            native_bp = self._cap_blob_parallelism(
+                min(effective, len(splits)), effective_bp)
+        native_batches = self._try_native_batches(
+            splits,
+            schema,
+            parallelism=effective,
+            blob_parallelism=native_bp,
+        )
+        if native_batches is not None:
+            return self._batches_to_arrow(native_batches, schema)
+
         if self._should_run_parallel(splits, effective):
             return self._to_arrow_parallel(splits, schema, effective, 
effective_bp)
 
-        batch_reader = self.to_arrow_batch_reader(splits, 
blob_parallelism=effective_bp)
+        batch_iterator = self._arrow_batch_generator(
+            splits, schema, effective_bp)
+        batch_reader = pyarrow.ipc.RecordBatchReader.from_batches(
+            schema, batch_iterator)
 
         table_list = []
         for batch in iter(batch_reader.read_next_batch, None):
             if batch.num_rows == 0:
                 continue
             table_list.append(self._try_to_pad_batch_by_schema(batch, schema))
 
-        if not table_list:
-            return pyarrow.Table.from_arrays([pyarrow.array([], 
type=field.type) for field in schema], schema=schema)
-        else:
-            return pyarrow.Table.from_batches(table_list)
+        return self._batches_to_arrow(table_list, schema)
+
+    @staticmethod
+    def _batches_to_arrow(batches, schema):
+        batches = [TableRead._try_to_pad_batch_by_schema(batch, schema)
+                   for batch in batches if batch.num_rows > 0]
+        if not batches:
+            return pyarrow.Table.from_arrays(
+                [pyarrow.array([], type=field.type) for field in schema],
+                schema=schema)
+        return pyarrow.Table.from_batches(batches)
+
+    def _try_native_batches(
+            self,
+            splits: List[Split],
+            schema: pyarrow.Schema,
+            parallelism: Optional[int] = None,
+            blob_parallelism: Optional[int] = None):
+        """Return Rust-read batches, or ``None`` when this read must fall 
back."""
+        if not self.table.options.native_read_enabled():
+            return None
+        # These Python-only output controls do not yet have native equivalents.
+        if self.include_row_kind or self.nested_name_paths:
+            return None
+        if self.table.options.file_format() not in _NATIVE_READ_FILE_FORMATS:
+            return None
+        if not splits:
+            return []
+        try:
+            from pypaimon.read.native_plan import (
+                native_blob_parallelism_available, native_read)
+        except Exception as e:
+            logger.warning(
+                "Native read failed, falling back to the Python reader: %s", e)
+            return None
+        blob_parallelism_available = native_blob_parallelism_available()
+        if ((blob_parallelism is not None and blob_parallelism > 1)
+                or self._deferred_blob_fields) and not 
blob_parallelism_available:
+            return None
+        rust_splits = []
+        for split in splits:
+            if isinstance(split, QueryAuthSplit):
+                return None
+            if not self._native_split_files_supported(
+                    split, blob_parallelism_available):
+                return None
+            rust_split = getattr(split, '_native_split', None)
+            if rust_split is None:
+                return None
+            rust_splits.append(rust_split)
+        native_bp = blob_parallelism if blob_parallelism_available else None
+        if (parallelism is not None
+                and self._should_run_parallel(splits, parallelism)):
+            try:
+                return self._native_batches_parallel(
+                    native_read, rust_splits, schema, parallelism, native_bp)
+            except _NativeReadSetupError as e:
+                logger.warning(
+                    "Native read failed, falling back to the Python reader: 
%s", e)
+                return None
+        try:
+            read_kwargs = {
+                'predicate': self.predicate,
+                'limit': self.limit,
+                'projection': [field.name for field in self.read_type],
+            }
+            if native_bp is not None:
+                read_kwargs['blob_parallelism'] = native_bp
+            batches = native_read(self.table, rust_splits, **read_kwargs)
+        except Exception as e:
+            logger.warning(
+                "Native read failed, falling back to the Python reader: %s", e)
+            return None
+        return self._convert_native_batches(batches, schema)
+
+    def _native_batches_parallel(
+            self, native_read, rust_splits, schema, effective,
+            blob_parallelism):
+        """Read contiguous split groups with independent Rust readers."""
+        workers = min(effective, len(rust_splits))
+        base_size, larger_groups = divmod(len(rust_splits), workers)
+        groups = []
+        offset = 0
+        for index in range(workers):
+            size = base_size + (1 if index < larger_groups else 0)
+            groups.append(rust_splits[offset:offset + size])
+            offset += size
+
+        remaining_state = _RemainingRows(self.limit)
+        results = [None] * len(groups)
+        with ThreadPoolExecutor(
+                max_workers=workers,
+                thread_name_prefix="pypaimon-native-read") as executor:
+            futures = {
+                executor.submit(
+                    self._read_native_split_group,
+                    native_read,
+                    group,
+                    schema,
+                    remaining_state,
+                    blob_parallelism,
+                ): index
+                for index, group in enumerate(groups)
+            }
+            for future in as_completed(futures):
+                results[futures[future]] = future.result()
+
+        return [batch for group_batches in results for batch in group_batches]
+
+    def _read_native_split_group(
+            self, native_read, rust_splits, schema, remaining_state,
+            blob_parallelism):
+        if remaining_state.exhausted():
+            return []
+        try:
+            read_kwargs = {
+                'predicate': self.predicate,
+                'limit': self.limit,
+                'projection': [field.name for field in self.read_type],
+            }
+            if blob_parallelism is not None:
+                read_kwargs['blob_parallelism'] = blob_parallelism
+            batches = native_read(self.table, rust_splits, **read_kwargs)
+        except Exception as e:
+            raise _NativeReadSetupError(str(e)) from e
+        result = []
+        for batch in batches:
+            if batch.num_rows == 0:
+                continue
+            allowed = remaining_state.try_consume(batch.num_rows)
+            if allowed == 0:
+                break
+            if allowed < batch.num_rows:
+                batch = batch.slice(0, allowed)
+            batch = self._project_batch_to_output(batch)
+            result.append(self._try_to_pad_batch_by_schema(batch, schema))
+            if remaining_state.exhausted():
+                break
+        return result
+
+    @staticmethod
+    def _native_split_files_supported(split, blob_parallelism_available=False):
+        for data_file in split.files:
+            file_name = data_file.file_name.lower()
+            if ('.vector.' not in file_name
+                    and not file_name.endswith(_NATIVE_READ_FILE_SUFFIXES)
+                    and not (blob_parallelism_available
+                             and 
file_name.endswith(_NATIVE_BLOB_FILE_SUFFIX))):
+                return False
+        return True
+
+    def _convert_native_batches(self, batches, schema):
+        """Apply PyPaimon's exact output limit lazily to native batches."""
+        remaining = self.limit
+        for batch in batches:
+            if batch.num_rows == 0:
+                continue
+            if remaining is not None and batch.num_rows > remaining:
+                batch = batch.slice(0, remaining)
+            batch = self._project_batch_to_output(batch)
+            yield self._try_to_pad_batch_by_schema(batch, schema)

Review Comment:
   Could we normalize native batches to the requested PyPaimon Arrow schema, 
including field types, before returning them on both the serial and parallel 
paths?
   
   For a BLOB column, PyPaimon expects `large_binary`, but the native reader 
returns `binary`. `_try_to_pad_batch_by_schema()` returns the batch unchanged 
when the column names match, so it does not handle this difference.
   
   I reproduced this with a two-row Data Evolution table containing an `INT` 
column and a BLOB column: the Python path returns `large_binary`, while native 
`to_arrow()` returns `binary`. More importantly, 
`to_arrow_batch_reader(...).read_all()` fails with `ArrowInvalid: Schema at 
index 0 was different`, because the reader declares `large_binary` but emits 
`binary` batches. Comparing only `to_pydict()` in the current tests misses this 
regression.
   
   Please add schema-equality assertions and a `read_all()` regression test, in 
addition to checking the row values.



##########
paimon-python/pypaimon/read/table_read.py:
##########
@@ -282,21 +304,199 @@ def to_arrow(
         if self.include_row_kind:
             schema = self._add_row_kind_to_schema(schema)
 
+        native_parallel = self._should_run_parallel(splits, effective)
+        native_bp = effective_bp
+        if native_parallel:
+            native_bp = self._cap_blob_parallelism(
+                min(effective, len(splits)), effective_bp)
+        native_batches = self._try_native_batches(
+            splits,
+            schema,
+            parallelism=effective,
+            blob_parallelism=native_bp,
+        )
+        if native_batches is not None:
+            return self._batches_to_arrow(native_batches, schema)
+
         if self._should_run_parallel(splits, effective):
             return self._to_arrow_parallel(splits, schema, effective, 
effective_bp)
 
-        batch_reader = self.to_arrow_batch_reader(splits, 
blob_parallelism=effective_bp)
+        batch_iterator = self._arrow_batch_generator(
+            splits, schema, effective_bp)
+        batch_reader = pyarrow.ipc.RecordBatchReader.from_batches(
+            schema, batch_iterator)
 
         table_list = []
         for batch in iter(batch_reader.read_next_batch, None):
             if batch.num_rows == 0:
                 continue
             table_list.append(self._try_to_pad_batch_by_schema(batch, schema))
 
-        if not table_list:
-            return pyarrow.Table.from_arrays([pyarrow.array([], 
type=field.type) for field in schema], schema=schema)
-        else:
-            return pyarrow.Table.from_batches(table_list)
+        return self._batches_to_arrow(table_list, schema)
+
+    @staticmethod
+    def _batches_to_arrow(batches, schema):
+        batches = [TableRead._try_to_pad_batch_by_schema(batch, schema)
+                   for batch in batches if batch.num_rows > 0]
+        if not batches:
+            return pyarrow.Table.from_arrays(
+                [pyarrow.array([], type=field.type) for field in schema],
+                schema=schema)
+        return pyarrow.Table.from_batches(batches)
+
+    def _try_native_batches(
+            self,
+            splits: List[Split],
+            schema: pyarrow.Schema,
+            parallelism: Optional[int] = None,
+            blob_parallelism: Optional[int] = None):
+        """Return Rust-read batches, or ``None`` when this read must fall 
back."""
+        if not self.table.options.native_read_enabled():
+            return None
+        # These Python-only output controls do not yet have native equivalents.
+        if self.include_row_kind or self.nested_name_paths:
+            return None
+        if self.table.options.file_format() not in _NATIVE_READ_FILE_FORMATS:
+            return None
+        if not splits:
+            return []
+        try:
+            from pypaimon.read.native_plan import (
+                native_blob_parallelism_available, native_read)
+        except Exception as e:
+            logger.warning(
+                "Native read failed, falling back to the Python reader: %s", e)
+            return None
+        blob_parallelism_available = native_blob_parallelism_available()
+        if ((blob_parallelism is not None and blob_parallelism > 1)

Review Comment:
   Support for `with_blob_parallelism()` does not guarantee that a pruning 
LIMIT is applied before BLOB payload resolution. Could we keep deferred-BLOB 
reads with a potentially pruning LIMIT on the Python path until the native 
reader has that capability, or pass the remaining row budget into Rust before 
payloads are fetched?
   
   The native iterator resolves a whole batch before 
`_convert_native_batches()` slices it to the requested limit. Making split 
processing serial only prevents cross-split prefetch; it does not prevent 
payload reads for discarded rows within that batch.
   
   With 32 rows referencing 32 distinct, valid local payload files, `LIMIT 1` 
returns the same one row on both paths, but filesystem open-event tracking 
shows that Python opens 1 payload file and native opens all 32. With two rows 
where only the second, discarded row references a missing payload, Python 
succeeds and native fails while opening that unused payload. This regresses the 
existing deferred-BLOB LIMIT behavior and can be expensive for image/video 
workloads.
   
   Please add a native-read regression test that checks payload I/O, not just 
the final row count.



##########
paimon-python/pypaimon/read/native_plan.py:
##########
@@ -327,6 +372,14 @@ def native_plan(
     # Trimmed primary keys decode per-file min/max keys (PK merge-on-read).
     kfields = table.trimmed_primary_keys_fields
     splits = [deserialize_split_v1(split.serialize(), pfields, kfields) for 
split in rust_splits]
+    if table.options.native_read_enabled():
+        # Retain the opaque Rust split next to the Python metadata view. The
+        # normal planner/reader contract remains a Python Split list, while
+        # native reads can consume the exact Rust split without a second lossy
+        # conversion. Any Python split transformation creates a fresh object
+        # without this marker and thus safely falls back to the Python reader.
+        for split, rust_split in zip(splits, rust_splits):
+            split._native_split = rust_split

Review Comment:
   The retained Rust split can become stale immediately after this assignment: 
`_restore_python_partition_paths()` updates the Python file paths in place, but 
does not update or invalidate `_native_split`.
   
   I reproduced this with a PyPaimon-written table partitioned by a string 
column whose value is `a/b`. The Python reader successfully reads the legacy 
`part=a/b/...` location, while enabling `read.native.enabled` makes Rust read 
`part=a%2Fb/...` and fail with `NotFound`. The existing legacy-partition 
integration test exercises native planning with Python reading, so it does not 
cover the new native-read path.
   
   Could we synchronize the corrected paths into the Rust split, or remove the 
native marker whenever path restoration changes the Python split so this case 
falls back safely?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to