leaves12138 commented on code in PR #9954:
URL: https://github.com/apache/paimon/pull/9954#discussion_r4046533769
##########
paimon-python/pypaimon/read/table_read.py:
##########
@@ -239,6 +262,10 @@ def _try_to_pad_batch_by_schema(batch:
pyarrow.RecordBatch, target_schema):
for field in target_schema:
if field.name in batch.schema.names:
col = batch.column(field.name)
+ if col.type != field.type:
Review Comment:
Fixing BLOB in paimon-rust#874 resolves that type mismatch, but valid
TIMESTAMP(0) and TIMESTAMP_LTZ(0) columns still fail this check.
PyPaimon's `PyarrowFieldParser` maps precision 0 to `timestamp[s]` (with UTC
for LTZ), while Rust's `timestamp_time_unit()` maps precisions 0 through 3 to
milliseconds. I reproduced this with an ordinary two-row Parquet table, using
the actual #874 binding: Python reading succeeds, but native reading raises
`TypeError: Batch field 'value' has type timestamp[ms], expected timestamp[s]`.
Both `to_arrow()` and `to_arrow_batch_reader(...).read_all()` fail; the LTZ
case fails with the corresponding UTC types.
Could we align the timestamp representations at the native/Python boundary,
or temporarily fall back for these types rather than only asserting equality?
Please cover both TIMESTAMP(0) and TIMESTAMP_LTZ(0) with real native-read
schema tests.
##########
paimon-python/pypaimon/read/table_read.py:
##########
@@ -400,24 +599,25 @@ def _cap_blob_parallelism(cls, workers: int,
blob_parallelism: int) -> int:
return max(1, cls._MAX_TOTAL_BLOB_WORKERS // workers)
def _should_run_parallel(
- self,
- splits: List[Split],
- effective: int,
+ self,
+ splits: List[Split],
+ effective: int,
) -> bool:
"""Decide whether to take the parallel read path.
``effective == 1`` falls back to the serial path (no thread pool
overhead, no behavior change). A single split is never
parallelized since there is nothing to fan out across.
"""
- deferred_limit_may_prune = (
- self.limit is not None
- and self._deferred_blob_fields
- and not self._limit_covers_all_splits(splits)
- )
+ deferred_limit_may_prune = self._deferred_blob_limit_may_prune(splits)
return (effective >= 2 and len(splits) >= 2
and not deferred_limit_may_prune)
+ def _deferred_blob_limit_may_prune(self, splits: List[Split]) -> bool:
+ return (self.limit is not None
+ and self._deferred_blob_fields
Review Comment:
Could we include descriptor-backed BLOB fields in the native LIMIT fallback
instead of relying only on `_deferred_blob_fields`?
`deferred_blob_field_names()` explicitly excludes fields configured through
`blob-descriptor-field`, so this guard is false for a Data Evolution table with
`blob-descriptor-field=payload`, even when LIMIT will discard rows. Rust then
resolves the whole batch before the Python wrapper applies the output limit.
I re-ran this against the current head and the binding from paimon-rust#874.
With 32 rows referencing 32 distinct, valid local payload files and LIMIT 1,
filesystem open-event tracking shows that Python opens 1 payload while native
opens all 32. This happens with both `to_arrow(..., parallelism=1)` and
`to_arrow_batch_reader(...).read_all()`. With two rows where only the second,
discarded descriptor references a missing payload, Python succeeds while native
fails with NotFound.
The added regression covers dedicated `.blob` fields, but not this
descriptor-backed case. Please extend the native-read eligibility check to
cover these payload-resolving fields, and add a real-binding regression for
`blob-descriptor-field` that checks discarded-payload I/O, not just the
returned row count.
--
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]