JingsongLi commented on code in PR #9954:
URL: https://github.com/apache/paimon/pull/9954#discussion_r4045965160
##########
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:
Fixed in be669d58e2. Native batches are now normalized against the full
requested Arrow schema: matching names are no longer sufficient, mismatched
column types are cast (including binary -> large_binary), and the batch is
rebuilt with the target schema. Added serial `RecordBatchReader.read_all()` and
parallel native-read regressions, and the live Data Evolution BLOB test now
asserts exact schema equality for both `to_arrow()` and `read_all()`.
##########
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:
Fixed in be669d58e2. A Data Evolution read with deferred BLOB fields now
falls back before native invocation whenever its LIMIT may prune rows. Limits
proven to cover all split rows remain native-eligible. The new live regression
tracks `BlobRef.to_data()` calls and verifies LIMIT 1 materializes exactly one
payload while `native_read` is never invoked.
--
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]