Skip to content

[python] Push down typed Variant projections - #10336

Open
XiaoHongbo-Hope wants to merge 15 commits into
apache:masterfrom
XiaoHongbo-Hope:codex/pypaimon-variant-projection
Open

XiaoHongbo-Hope wants to merge 15 commits into
apache:masterfrom
XiaoHongbo-Hope:codex/pypaimon-variant-projection

Conversation

@XiaoHongbo-Hope

@XiaoHongbo-Hope XiaoHongbo-Hope commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor

Purpose

Push selected VARIANT fields into the native reader without exposing reader-specific struct children to callers. Depends on apache/paimon-rust#1009.

API

builder.with_projection({
    "id": "id",
    "x": "try_variant_get(payload, '$.x', 'float')",
    "y": "variant_get(payload, '$.y', 'float')",
})

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_get returns NULL on cast failure, while variant_get raises an error. Outputs are flat named columns. PyPaimon derives a Paimon ROW read type and passes it to Rust with_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

  • 294 focused tests passed, including real local-file native batch/stream reads, output schema, SQL NULL and both cast policies.
  • Added a named-projection test to the Python 3.6/3.7 CI subset. Flake8 and git diff --check passed; 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.
  • An earlier same-snapshot A/B of the same Rust read-type path (using the previous public API) measured data-supply throughput of 270.15 → 343.97 samples/s and Python-visible topic Arrow data of 426.4 → 1.21 MB. The new mapping API has not been separately benchmarked; these are not model-training throughput results.

@XiaoHongbo-Hope
XiaoHongbo-Hope marked this pull request as ready for review October 3, 2026 14:58
Comment thread docs/docs/pypaimon/data-types.md Outdated
'paths': ["$['state.x']", "$['action.y']"],
'target_type': pa.float32(),
},
},

@XiaoHongbo-Hope XiaoHongbo-Hope Oct 4, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@JingsongLi Could you please let me know your suggestion about about the user API in pypaimon side?

@JingsongLi

Copy link
Copy Markdown
Contributor

The extraction pushdown is useful, but I would prefer to express it through projection expressions rather than expose variant_fields as a public API. Currently, callers have to describe both the output projection and the native reader's extraction plan, and the result exposes positional struct children ("0", "1", ...).

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 List[str] form could retain its column-name semantics, while the mapping form would describe named projection expressions. The output would be id, x, y, with each extraction specifying its own path, type, and error behavior. try_variant_get matches the current default fail_on_error=False; variant_get would express strict behavior.

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 read_type, pass it through with_read_type, and project the internal struct children into the requested output columns. This also follows Spark's Variant extraction pushdown model: users write logical expressions, and the planner derives the reader request.

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.

@XiaoHongbo-Hope

Copy link
Copy Markdown
Contributor Author

The extraction pushdown is useful, but I would prefer to express it through projection expressions rather than expose variant_fields as a public API. Currently, callers have to describe both the output projection and the native reader's extraction plan, and the result exposes positional struct children ("0", "1", ...).

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 List[str] form could retain its column-name semantics, while the mapping form would describe named projection expressions. The output would be id, x, y, with each extraction specifying its own path, type, and error behavior. try_variant_get matches the current default fail_on_error=False; variant_get would express strict behavior.

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 read_type, pass it through with_read_type, and project the internal struct children into the requested output columns. This also follows Spark's Variant extraction pushdown model: users write logical expressions, and the planner derives the reader request.

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]],

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

API Looks good to me.

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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 JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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 $.a'b, but this parser records $.ab. With real persisted data containing both keys, the current mapping returns [9.5, NULL] instead of [1.25, 3.5]; a native read-type control with the decoded SQL path returns the correct values. This reproduces for append and DE tables in batch and stream reads, so it can silently read another field rather than just reject unusual syntax. Decode SQL string literals according to the projection API contract (or explicitly reject unsupported escaping), and add an apostrophe-key regression. The doubled-quote convention is documented in SQL lexical syntax.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants