leaves12138 commented on code in PR #10320:
URL: https://github.com/apache/paimon/pull/10320#discussion_r4134779549
##########
paimon-python/pypaimon/data/variant_path.py:
##########
@@ -1646,6 +1780,106 @@ def variant_get(column, path, data_type=None):
return _variant_get(column, {path: data_type})[path]
+@_with_metadata_cache
+def variant_to_pylist(column, fields: Sequence[str]):
+ """Decode selected top-level VARIANT object fields to Python values.
+
+ Values use their natural Python types, even when a field has different
+ types in different rows. A missing field is omitted from its row's dict;
+ a VARIANT null is present with value None, and a SQL null row returns None.
+ Field names are literal names, not VARIANT path expressions.
+ """
+ if isinstance(fields, (str, bytes)):
+ raise TypeError("VARIANT fields must be a sequence of field names")
+ try:
+ names = tuple(dict.fromkeys(fields))
+ except (TypeError, ValueError):
+ raise TypeError(
+ "VARIANT fields must be a sequence of field names") from None
+ if any(not isinstance(name, str) for name in names):
+ raise TypeError("VARIANT field names must be strings")
+
+ chunks, _, _ = _variant_chunks(column)
+ result = []
+ slot_cache = LRUCache(maxsize=64)
+ for chunk in chunks:
+ values = _BinaryValues(chunk.field(0))
+ metadata = _BinaryValues(chunk.field(1))
+ valid = set(map(int, _valid_row_indices(
+ chunk, values, chunk.field(1))))
+ for row in range(len(chunk)):
+ if row not in valid:
+ result.append(None)
+ continue
+ value = bytes(values.view(row))
+ row_metadata = bytes(metadata.view(row))
+ if not value or (value[0] & 0x3) != _OBJECT:
+ raise TypeError("VARIANT root must be an object")
+ variant = GenericVariant(value, row_metadata)
+ key_ids = _cached_metadata_key_ids(row_metadata)
+ targets = {key_ids[name]: name for name in names
+ if name in key_ids}
+ (size, id_width, id_start, offset_start, data_start,
+ offsets, ordered_offsets) = _flat_object_layout(value)
+ if data_start + int(offsets[-1]) != len(value):
+ _malformed("trailing bytes after root value")
+ layout = (row_metadata, value[id_start:offset_start],
+ id_width, tuple(sorted(targets)))
+ selected_slots = slot_cache.get(layout)
+ if selected_slots is None:
+ field_ids = _flat_object_unsigneds(
+ value, id_start, size, id_width)
+ if len(np.unique(field_ids)) != size:
+ _malformed("duplicate object field id")
+ selected_slots = tuple(
+ (slot, targets[int(key_id)])
+ for slot, key_id in enumerate(field_ids)
+ if int(key_id) in targets)
+ slot_cache[layout] = selected_slots
+ if not targets:
+ result.append({})
+ continue
+ slots = np.fromiter(
+ (slot for slot, _ in selected_slots),
+ dtype=np.int64, count=len(selected_slots))
+ starts = offsets[slots]
+ if np.any(starts >= offsets[-1]):
+ _malformed("invalid object offset")
+ next_indices = np.searchsorted(
+ ordered_offsets, starts, side='right')
+ ends = ordered_offsets[next_indices]
+ headers = np.frombuffer(value, dtype=np.uint8)[
+ data_start + starts]
Review Comment:
[P1] Widen offsets before adding the object header size.
`starts` retains the encoded offset dtype (`uint8` or `uint16`), so
`data_start + starts` is performed in that narrow dtype. Valid offsets can wrap
when converted from data-relative to buffer-relative positions; on NumPy 2, a
larger Python scalar can instead raise OverflowError. This makes the new
boundary check reject valid objects.
Minimal reproduction on this head:
```python
from pypaimon.data.generic_variant import GenericVariant
from pypaimon.data import variant_to_pylist
obj = GenericVariant.from_python({
'field%04d' % index: None for index in range(100)
})
column = GenericVariant.to_arrow_array([obj])
variant_to_pylist(column, ['field0099'])
```
Expected `[{'field0099': None}]`; actual `ValueError: MALFORMED_VARIANT:
child size does not match container offsets`. Here `data_start=203`, and the
last relative offset 99 wraps instead of indexing byte 302. A valid
128-null-field object raises `OverflowError: Python integer 259 out of bounds
for uint8`, and selecting the last field of a valid 5,362-double-field object
also fails with uint16 offsets.
Please promote the offsets to an index-sized integer before the addition,
and cover both one-byte/two-byte offset boundary crossings. Widening the
offsets in a local diagnostic makes all three reproductions pass.
##########
paimon-python/pypaimon/data/variant_path.py:
##########
@@ -1579,6 +1584,133 @@ def _rowwise_replace_chunk(
chunk, values, data, data_start, rebuilt_rows)
+def _flat_object_unsigneds(value, start, count, width):
+ """Read an object header array without creating Python ints per field."""
+ if width in (1, 2, 4):
+ return np.frombuffer(
+ value, dtype=np.dtype('<u%d' % width), count=count,
+ offset=start)
+ raw = np.frombuffer(
+ value, dtype=np.uint8, count=count * 3,
+ offset=start).reshape(count, 3)
+ return (raw[:, 0].astype(np.uint32)
+ | (raw[:, 1].astype(np.uint32) << 8)
+ | (raw[:, 2].astype(np.uint32) << 16))
+
+
+def _flat_object_layout(value):
+ """Validate a root object and read its offsets as a NumPy array."""
+ limit = len(value)
+ _require_range(0, 2, limit)
+ type_info = (value[0] >> 2) & 0x3F
+ size_width = _U32_SIZE if ((type_info >> 4) & 0x1) else 1
+ _require_range(1, size_width, limit)
+ size = _read_unsigned(value, 1, size_width)
+ id_width = ((type_info >> 2) & 0x3) + 1
+ offset_width = (type_info & 0x3) + 1
+ id_start = 1 + size_width
+ offset_start = id_start + size * id_width
+ data_start = offset_start + (size + 1) * offset_width
+ _require_range(0, data_start, limit)
+ offsets = _flat_object_unsigneds(
+ value, offset_start, size + 1, offset_width)
+ end_offset = int(offsets[-1])
+ if size:
+ ordered = np.sort(offsets)
+ if (int(ordered[0]) != 0 or int(ordered[-1]) != end_offset
+ or np.any(np.diff(ordered) == 0)):
+ _malformed("invalid object offsets")
+ else:
+ if end_offset != 0:
+ _malformed("invalid object offsets")
+ ordered = offsets
+ _require_range(data_start, end_offset, limit)
+ return size, id_width, id_start, offset_start, data_start, offsets, ordered
+
+
+def _flat_object_get_chunk(chunk, values, parsed):
+ """Look up wide, top-level object fields together instead of per path."""
+ if (len(parsed) < 2
+ or any(len(path) != 1 or path[0][0] != 'key'
+ for _, path, _ in parsed)):
+ return None
+
+ valid_rows = _valid_row_indices(chunk, values, chunk.field(1))
+ if not len(valid_rows):
+ return [pa.nulls(len(chunk), type=data_type)
+ for _, _, data_type in parsed]
+ first_value = values.view(int(valid_rows[0]))
+ if not first_value or (first_value[0] & 0x3) != _OBJECT:
+ return None
+ first_size = _checked_object_layout(
+ first_value, 0, len(first_value))[0]
+ # Keep the vectorized reader for narrow objects or a few requested keys.
+ if first_size * len(parsed) < 1024:
Review Comment:
[P2] Preserve the existing batch-vectorized path for homogeneous float
objects.
This guard only uses field-count times path-count, so it also redirects
relatively small objects with many same-layout rows from the existing NumPy
batch reader into per-row layout validation and decoding. That introduces a
performance regression in the existing variant_get API, without applications
opting into the new API.
I compared this head directly with base d3393ebf on the same prebuilt Arrow
arrays, using exact Float64 paths, identical object layouts but varying row
values, a warm-up, and the median of five alternating calls. Outputs are
Arrow-equal:
- 4,096 rows x 128 Float64 fields, selecting 16 top-level fields: base 0.137
s, head 0.499 s (~3.65x slower).
- 1,024 rows x 128 Float64 fields, selecting 8: base 0.036 s, head 0.081 s
(~2.26x slower).
The reported gain for much wider objects can still hold, but the current
threshold also captures these batches. Please retain the homogeneous
floating-point vectorized route when it is cheaper, or batch the shared lookup
rather than decoding every row in Python. Include a same-layout multi-row
comparison alongside the very-wide-object benchmark so this dispatch regression
is covered.
--
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]