JingsongLi commented on code in PR #10336:
URL: https://github.com/apache/paimon/pull/10336#discussion_r4180334802
##########
paimon-python/pypaimon/read/read_builder.py:
##########
@@ -186,6 +208,78 @@ def read_type(self) -> List[DataField]:
# Helpers
# ------------------------------------------------------------------
+ def _parse_expression_projection(self, expressions: Dict[str, str]):
+ if not expressions:
+ raise ValueError("Projection expression mapping must not be empty")
+ table_fields = self.table.fields
+ if self.table.options.row_tracking_enabled():
+ table_fields =
SpecialFields.row_type_with_row_tracking(table_fields)
+ field_map = {field.name: field for field in table_fields}
+ projection = []
+ variants = {}
+ outputs = []
+ direct_columns = set()
+ for alias, expression in expressions.items():
+ if not isinstance(alias, str) or not alias:
+ raise TypeError("Projection output names must be non-empty
strings")
+ if alias == ROW_KIND_COLUMN:
+ raise ValueError("Projection output name %r is reserved" %
alias)
+ if not isinstance(expression, str) or not expression:
+ raise TypeError("Projection expressions must be non-empty
strings")
+ if expression in field_map:
+ source, child = expression, None
+ direct_columns.add(source)
+ else:
+ try:
+ call = ast.parse(expression, mode='eval').body
+ except SyntaxError as error:
+ raise ValueError(
+ "Unsupported projection expression %r" % expression
+ ) from error
+ if (not isinstance(call, ast.Call)
+ or not isinstance(call.func, ast.Name)
+ or call.func.id.lower() not in ('variant_get',
'try_variant_get')
+ or len(call.args) != 3 or call.keywords
+ or not (isinstance(call.args[0], ast.Name)
+ or _string_literal(call.args[0]) is not None)
+ or any(_string_literal(arg) is None
+ for arg in call.args[1:])):
+ raise ValueError(
+ "Unsupported projection expression %r" % expression)
+ source = (call.args[0].id
+ if isinstance(call.args[0], ast.Name)
+ else _string_literal(call.args[0]))
+ field = field_map.get(source)
+ if (field is None
+ or not isinstance(field.type, AtomicType)
+ or field.type.type.upper() != 'VARIANT'):
+ raise ValueError(
+ "Variant extraction requires a VARIANT column: %r"
+ % source)
+ path, target_type = (_string_literal(arg) for arg in
call.args[1:])
+ if ';' in path:
+ raise ValueError(
+ "Variant extraction path must not contain ';': %s"
+ % path)
+ if target_type.lower() != 'float':
+ raise ValueError(
+ "Only float32 Variant extractions are supported")
+ options = variants.setdefault(source, {
+ 'paths': [], 'target_type': pyarrow.float32(),
+ 'fail_on_error': [],
+ })
+ child = len(options['paths'])
+ options['paths'].append(path)
+ options['fail_on_error'].append(
+ call.func.id.lower() == 'variant_get')
+ if source not in projection:
+ projection.append(source)
+ outputs.append((alias, source, child))
Review 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.
##########
paimon-python/pypaimon/read/read_builder.py:
##########
@@ -186,6 +208,78 @@ def read_type(self) -> List[DataField]:
# Helpers
# ------------------------------------------------------------------
+ def _parse_expression_projection(self, expressions: Dict[str, str]):
+ if not expressions:
+ raise ValueError("Projection expression mapping must not be empty")
+ table_fields = self.table.fields
+ if self.table.options.row_tracking_enabled():
+ table_fields =
SpecialFields.row_type_with_row_tracking(table_fields)
+ field_map = {field.name: field for field in table_fields}
+ projection = []
+ variants = {}
+ outputs = []
+ direct_columns = set()
+ for alias, expression in expressions.items():
+ if not isinstance(alias, str) or not alias:
+ raise TypeError("Projection output names must be non-empty
strings")
+ if alias == ROW_KIND_COLUMN:
+ raise ValueError("Projection output name %r is reserved" %
alias)
+ if not isinstance(expression, str) or not expression:
+ raise TypeError("Projection expressions must be non-empty
strings")
+ if expression in field_map:
+ source, child = expression, None
+ direct_columns.add(source)
+ else:
+ try:
+ call = ast.parse(expression, mode='eval').body
Review 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 $.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](https://www.postgresql.org/docs/15/sql-syntax-lexical.html#SQL-SYNTAX-CONSTANTS).
--
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]