JingsongLi commented on code in PR #9954:
URL: https://github.com/apache/paimon/pull/9954#discussion_r4046260372
##########
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:
Follow-up: fixed the root cause in apache/paimon-rust#874. BLOB now maps to
Arrow LargeBinary across scalar/nested format reads, descriptor resolution,
Data Evolution, and DataFusion, while BINARY/VARBINARY remain Binary. This PR
no longer casts legacy Binary BLOB batches or probes older unpublished Rust
bindings; it requires matching native field types and keeps the serial read_all
plus parallel schema regressions.
--
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]