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]

Reply via email to