Repository navigation
Conversation
|
Thanks @YutaLin. LGTM overall. Could you address feedback, then I'll kick off CI |
|
Hi @andygrove, thanks for the review! About "Spark accepts utf8 as an alias for UTF-8", spark only supports alias before 3.5, because it uses JDK https://spark.apache.org/docs/4.0.0/sql-migration-guide.html#upgrading-from-spark-sql-35-to-40
|
…atafusion-comet into 3183_support_spark_expression_encode
…expression_encode # Conflicts: # docs/source/user-guide/latest/expressions.md
coderfender
left a comment
There was a problem hiding this comment.
Left some minor comments but overall looks good @YutaLin
|
Hi @coderfender |
|
Hi @coderfender |
|
Sorry for the delay. Let me look into it shortly @YutaLin |
|
Thanks @YutaLin. This looks good. Could you fix the conflicts? |
…expression_encode # Conflicts: # docs/source/contributor-guide/spark_expressions_support.md # docs/source/user-guide/latest/expressions.md # spark/src/main/spark-4.0/org/apache/comet/shims/CometExprShim.scala # spark/src/main/spark-4.1/org/apache/comet/shims/CometExprShim.scala # spark/src/main/spark-4.2/org/apache/comet/shims/CometExprShim.scala
…atafusion-comet into 3183_support_spark_expression_encode
|
Hi @andygrove, |
|
Seems like there are still some more conflicts @YutaLin . |
…expression_encode # Conflicts: # docs/source/user-guide/latest/expressions.md
|
Hi @coderfender Could you please take another look when you have a chance? |
|
Sure. Yeah have been working on some doc changes causing merge conflicts. I will review it shortly |
andygrove
left a comment
There was a problem hiding this comment.
Nice approach lowering encode(col, 'utf-8') to a binary cast instead of writing a native implementation. Most of the earlier feedback looks addressed (null charset guard, the utf8 alias discussion, the trait rename, the 4.x consolidation), so I only have one substantive point plus a couple of minor ones.
Invalid UTF-8 input. The equivalence "encode(str, 'utf-8') equals cast(str AS binary)" only holds when the input string is valid UTF-8. Spark's Encode.encode guards the raw-bytes fast path with input.isValid:
if ("UTF-8".equalsIgnoreCase(toCharset) && input.isValid) return input.getBytesWhen a UTF8String holds malformed UTF-8, Spark does not return the raw bytes. It routes through a CharsetEncoder and, with the default non-legacy error action, either substitutes or throws malformedCharacterCoding. Spark 3.x reaches a similar outcome via input.toString.getBytes(charset), which turns invalid sequences into the replacement character first.
The cast approach here always returns the raw underlying bytes, so for a valid-UTF-8 column the two match perfectly, but for an invalid-UTF-8 string Comet would silently produce a different result than Spark's encode, with no fallback since the charset is still utf-8. Invalid UTF-8 in a StringType column is uncommon but reachable, for example via cast(binary_col as string) over non-UTF-8 bytes. Since expressions.md marks encode as fully compatible, could we either note this as a known limitation or add a test that feeds an invalid-UTF-8 string through both engines to confirm the divergence? At minimum a comment next to the "byte-equivalent to cast" note pointing out the isValid caveat would help the next reader.
StaticInvoke match robustness (minor). The 4.x match compares arguments.size == 4 and the full inputTypes sequence including StringTypeWithCollation(supportsTrimCollation = true). If a future 4.x minor version tweaks that shape, the match silently fails and encode falls back. That is safe but easy to miss. A looser match keyed on staticObject == classOf[Encode] and functionName == "encode" with positional argument extraction might be more robust across versions, if that is acceptable.
Test file (minor). encode.sql is missing a trailing newline, which some lint or RAT steps dislike. An invalid-UTF-8 input case would also directly probe the compatibility question above.
Overall this looks close. The invalid-UTF-8 behavior is really the only item I would want to settle before merge.
|
is this PR related to apache/datafusion#21331 ? |
…expression_encode # Conflicts: # spark/src/main/scala/org/apache/comet/serde/strings.scala
|
Hi @andygrove |
Lowering Three things. The malformed-UTF-8 divergence is not reported to users Spark replaces malformed bytes during
If you would rather not gate it, that is a discussion worth having explicitly, but it should not be an implicit consequence of where the code happens to live. Charset alias matching str.toString.toLowerCase(Locale.ROOT) == "utf-8"Java's Would Duplication between the 3.4 and 3.5 shims The additions to |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
encodefell back to Spark. This PR enables native execution for literal UTF-8 charsets. - Design approach: Reuse
CAST(string AS binary)through a shared helper, with version-specific expression dispatch. - Correctness / compatibility analysis: Valid UTF-8, nulls and empty strings match Spark. The previously reported malformed-input P2 concern remains for raw strings arriving through the JVM/Arrow handoff: Spark encodes byte
FFasEFBFBD, while the new lowering preservesFF. Documentation now acknowledges this, but there is no compatibility gate. This duplicates existing feedback, so it is not repeated as a new finding. - Key design decisions: Match UTF-8 case-insensitively and retain Spark fallback for other charsets, aliases and dynamic charset expressions. Shared Spark 4.x dispatch avoids repeated implementations.
- Implementation sketch: Spark 3.x matches
Encode; Spark 4.x matches itsStaticInvokereplacement. Both serialize the value into the existing binary-cast path. - Behavioral changes worth calling out: Eligible projections now execute natively. Reusing the existing cast avoids adding a separate codec implementation. No reproducible performance regression or unnecessary abstraction requiring P1/P2 action was found; performance was not benchmarked.
- Suggested improvements: Resolve the existing malformed-input concern through Spark-equivalent handling or compatibility fallback. No additional, nonduplicate introduced P1/P2 issues found within this review.
Reviewed all eight changed files against base 0761e549a43fee618af462a305645e70d663a380 at full head 8f855d873f2bb0bf291947b714f82d333fdf1b14. The PR remains non-draft. Read existing reviews, issue comments, inline comments and threads.
Routed skills: review-comet-pr, review-comet-expression-pr and audit-comet-expression. This head uses .claude/skills; the newer .ai router and expression sibling were read from another local Comet checkout.
Exact-head CI: GitHub reports zero check runs, workflow runs and commit statuses. The aggregate status is pending, not a passing result.
Validation: Compared Spark implementations and tests for 3.4.3, 3.5.8, 4.0.2, 4.1.2 and 4.2.0, plus the native/Arrow cast path. Spark 4.1.3 reference checks passed for ANSI modes, nulls, Unicode, charset casing, aliases and malformed-byte behavior. git diff --check passed. Comet integration tests could not run: the offline native build lacked async-compression v0.4.42, and Maven bootstrap failed resolving repo.maven.apache.org. No end-to-end Comet execution was validated. Project files remain unchanged.
|
#3183 was closed by #5037. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: At the specified base,
encodefalls back to Spark. This PR enables native execution for literal UTF-8 charsets. - Design approach: Lower eligible calls to the existing
CAST(string AS binary)implementation, avoiding a separate codec kernel. - Correctness / compatibility analysis: Valid UTF-8, nulls and empty strings match Spark. The previously reported malformed-input P2 remains unresolved: for raw malformed strings entering through the JVM/Arrow boundary, Spark encodes byte
FFasEFBFBD, while the new cast preservesFF. Documentation acknowledges this, but the lowering bypasses compatibility gating. This is covered by existing feedback, so it is not repeated as a new finding. - Key design decisions: Match literal
utf-8case-insensitively and retain fallback for other charsets, aliases and dynamic charset expressions. This preserves Spark 4.x charset restrictions. - Implementation sketch: Spark 3.x matches
Encode; shared Spark 4.x handling matches itsStaticInvokereplacement. Both serialize the value throughCometCast.castToProto. - Behavioral changes worth calling out: Eligible projections now execute natively. Reusing the existing cast keeps implementation complexity and additional runtime work small. No reproducible performance regression or P1/P2 abstraction concern was identified. Performance was not benchmarked.
- Suggested improvements: Resolve the existing malformed-input concern with Spark-equivalent handling or compatibility fallback. No additional, nonduplicate introduced P1/P2 issues found within this review.
Reviewed all eight changed files against base 0761e549a43fee618af462a305645e70d663a380 at full head 8f855d873f2bb0bf291947b714f82d333fdf1b14. The PR remains non-draft. Existing reviews, issue comments, inline comments and threads were read, excluding Copilot.
Routed skills: review-comet-pr, review-comet-expression-pr and audit-comet-expression. This checkout uses .claude/skills; the newer .ai router and expression sibling were read from another local Comet checkout.
Exact-head CI: Zero check runs, workflow runs and commit statuses. Aggregate status is pending, not passing.
Validation: Compared Spark sources and tests for 3.4.3, 3.5.8, 4.0.2, 4.1.2 and 4.2.0, plus the native/Arrow cast path. Spark 3.5.9 reference checks passed for Unicode, nulls, empty strings, charset casing, aliases, ANSI modes and malformed-byte behavior. git diff --check passed. Comet integration tests could not run: the offline native build lacked aws-config v1.10.0, and Maven compilation was blocked by DNS resolution of repo.maven.apache.org. No end-to-end Comet execution was validated. Project files and GitHub state remain unchanged.
Which issue does this PR close?
Closes #3183
Rationale for this change
Support expression Encode
What changes are included in this PR?
How are these changes tested?
Add encode.sql and run it in spark 3.4/3.5/4.0