[python] Push down typed Variant projections - #10336
XiaoHongbo-Hope wants to merge 15 commits into
Conversation
| 'paths': ["$['state.x']", "$['action.y']"], | ||
| 'target_type': pa.float32(), | ||
| }, | ||
| }, |
There was a problem hiding this comment.
@JingsongLi Could you please let me know your suggestion about about the user API in pypaimon side?
|
The extraction pushdown is useful, but I would prefer to express it through projection expressions rather than expose Could we follow the alias-to-SQL-expression pattern used by Lance's scanner and LanceDB's select? For example: builder.with_projection({
"id": "id",
"x": "try_variant_get(payload, '$.x', 'float')",
"y": "try_variant_get(payload, '$.y', 'float')",
})The existing Paimon Rust already has the expression-to-read-type rewrite in #460. We could extract/reuse that logic to collect the Variant extractions, generate the existing To keep this PR focused, the first implementation could support only ordinary columns and float32 Variant extractions, explicitly rejecting other expressions for now. It should preserve SQL null/cast/error semantics and verify both the named output schema and the actual extraction read type. The Lance comparison here is about the public API shape; it does not assume equivalent field-level I/O pushdown for Lance JSON. |
Thanks, updated the PR to use named projection expressions. It follows #460’s read-type pushdown path, without reusing its DataFusion-specific optimizer. |
| preserve the pre-existing contract. | ||
| def with_projection( | ||
| self, | ||
| projection: Union[List[str], Dict[str, str]], |
JingsongLi
left a comment
There was a problem hiding this comment.
I reproduced two correctness issues in the updated named projection API; details and repair suggestions are inline. The focused validation passed 161 Python tests, 66 subtests, and 10 Java tests, including native projection reads, but the two additional reproductions expose cases not covered by those tests.
| direct_columns.add(source) | ||
| else: | ||
| try: | ||
| call = ast.parse(expression, mode='eval').body |
There was a problem hiding this comment.
[P2] Preserve SQL string-literal semantics when parsing paths
ast.parse() applies Python literal rules, so SQL's doubled single quote is silently removed by concatenating adjacent Python string literals. For example, try_variant_get(payload, '$["it''s"]', 'float') is parsed as the path $["its"].
I reproduced this with an actual table containing {"it's": 1.0, "its": 2.0} and the native reader: with_projection({'x': expression}) returns 2.0, while SQLContext executing the exact same expression returns 1.0. Both variant_get and try_variant_get are affected. This reads a different field without reporting an error.
Please parse literals using SQL rules, or at least reject this syntax explicitly instead of silently rewriting the path. A regression test should compare the named projection with SQL for a key containing an apostrophe.
| effective_bp = self._resolve_blob_parallelism(blob_parallelism) | ||
| effective = self._effective_parallelism(parallelism, len(splits)) | ||
| schema = PyarrowFieldParser.from_paimon_schema(self.read_type) | ||
| schema = self._output_arrow_schema() |
There was a problem hiding this comment.
[P2] Apply named projection after converting merged rows to Arrow
The output schema can have more fields than the physical read type because multiple aliases may reference the same source column, and _parse_expression_projection() deduplicates those source columns. The Python row-reader branch in _arrow_batch_generator() converts physical row tuples directly with this output schema and never applies _project_batch_to_output(); the parallel row-reader branch has the same issue.
I reproduced this on a primary-key table with native reading disabled, after two overlapping commits (id=1, val=20, then id=1, val=30):
builder = table.new_read_builder().with_projection({
'one': 'id',
'two': 'id',
'value': 'val',
})The physical read type is ['id', 'val'], but the output schema is ['one', 'two', 'value']. Both to_arrow(splits, parallelism=1) and to_arrow_batch_reader(splits, parallelism=1) raise KeyError: 'value' instead of returning one=1, two=1, value=30.
Please convert merged rows with the physical schema first, then apply the existing expression projection, preserving row kind. A regression test should cover repeated source columns on a non-raw-convertible primary-key split.
JingsongLi
left a comment
There was a problem hiding this comment.
Re-reviewed the named projection API at the current head. Requirement fit: SUPPORTED; deriving the internal extraction read type from named expressions addresses the API concern and keeps the Rust pushdown useful. Implementation: two P2 findings below.
Validation: 273 focused tests passed on Python 3.11 (two native cases skipped); with the paired Rust extension on its Python 3.13 runtime, 253 tests and 66 subtests passed. Six additional real append/DE batch/stream and strict-cast checks passed. Boundary probes reproduced the repeated-column failure on a PK table while append controls worked, and reproduced the SQL-literal retargeting in all four append/DE × batch/stream combinations. Flake8 and diff checks passed. Ray was not installed locally; affected Python CI lanes are green. Deployment still requires the compatible Rust #1009 extension.
| call.func.id.lower() == 'variant_get') | ||
| if source not in projection: | ||
| projection.append(source) | ||
| outputs.append((alias, source, child)) |
There was a problem hiding this comment.
[P2] Apply named projections after converting the physical row tuple
A mapping can legally project one ordinary column more than once, e.g. {'x': 'id', 'copy': 'id'}. These lines deduplicate the physical read columns but preserve both output entries. In the Python PK row-reader path, _arrow_batch_generator passes the final alias schema directly to _convert_rows_to_arrow_batches_with_row_kind without applying _project_batch_to_output, so the shorter physical tuple cannot fill that schema. On a real bucket=1 PK table with two overlapping commits, both to_arrow() and to_arrow_batch_reader() raise KeyError('copy'); adding a projected v column raises KeyError('v'). The same queries on an append table return both aliases correctly. Convert row tuples using the physical source schema first, then run the shared expression projection before emitting the final batch, and cover duplicate ordinary columns on a merged PK read.
| direct_columns.add(source) | ||
| else: | ||
| try: | ||
| call = ast.parse(expression, mode='eval').body |
There was a problem hiding this comment.
[P2] Preserve SQL escaping in projection string literals
Using Python ast.parse for the requested SQL expression syntax silently changes SQL doubled quotes into Python adjacent-string concatenation. For example, try_variant_get(payload, '$.a''b', 'float') should request the supported JSON path
Purpose
Push selected VARIANT fields into the native reader without exposing reader-specific struct children to callers. Depends on apache/paimon-rust#1009.
API
The existing
List[str]form is unchanged. The mapping form currently supports ordinary columns and FLOAT32 VARIANT extraction only; other expressions fail explicitly.try_variant_getreturns NULL on cast failure, whilevariant_getraises an error. Outputs are flat named columns. PyPaimon derives a Paimon ROW read type and passes it to Rustwith_read_type; batch, stream, and Ray share this path. VARIANT extraction requires native read and does not silently fall back. Row iterators and Torch row mode explicitly reject the mapping.Validation
git diff --checkpassed; Python 3.6 parsed the actual string-literal helper correctly. A Ray integration test covers named output and an empty table; Ray is not installed locally, so that test awaits CI.