Skip to content

feat: Support Spark Expression Encode - #4315

Open
YutaLin wants to merge 23 commits into
apache:mainfrom
YutaLin:3183_support_spark_expression_encode
Open

YutaLin wants to merge 23 commits into
apache:mainfrom
YutaLin:3183_support_spark_expression_encode

Conversation

@YutaLin

@YutaLin YutaLin commented May 13, 2026

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #3183

Rationale for this change

Support expression Encode

What changes are included in this PR?

  • Add StringEncode in string serde
  • Update shims in spark3.4/3.5/4.0/4.1/4.2 to catch Encode

How are these changes tested?

Add encode.sql and run it in spark 3.4/3.5/4.0

Comment thread spark/src/main/spark-4.0/org/apache/comet/shims/CometExprShim.scala Outdated
Comment thread spark/src/main/scala/org/apache/comet/serde/strings.scala Outdated
@andygrove

Copy link
Copy Markdown
Member

Thanks @YutaLin. LGTM overall. Could you address feedback, then I'll kick off CI

@YutaLin

YutaLin commented May 13, 2026

Copy link
Copy Markdown
Contributor Author

Hi @andygrove, thanks for the review!
I've extract encode method and add null check.

About "Spark accepts utf8 as an alias for UTF-8", spark only supports alias before 3.5, because it uses JDK Charset.forName. After 4.0, it has a whitelist check, so it doesn't support alias. I'd suggest we keep only utf-8 now, WDYT?

https://spark.apache.org/docs/4.0.0/sql-migration-guide.html#upgrading-from-spark-sql-35-to-40

Since Spark 4.0, the encode() and decode() functions support only the following charsets ‘US-ASCII’, ‘ISO-8859-1’, ‘UTF-8’, ‘UTF-16BE’, ‘UTF-16LE’, ‘UTF-16’, ‘UTF-32’. To restore the previous behavior when the function accepts charsets of the current JDK used by Spark, set spark.sql.legacy.javaCharsets to true.

@YutaLin
YutaLin requested a review from andygrove May 13, 2026 22:17

@coderfender coderfender 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.

Left some minor comments but overall looks good @YutaLin

Comment thread spark/src/main/scala/org/apache/comet/serde/strings.scala
Comment thread spark/src/main/spark-4.x/org/apache/comet/shims/ShimCometExprs.scala Outdated
Comment thread spark/src/test/resources/sql-tests/expressions/string/encode.sql
Comment thread spark/src/main/spark-4.1/org/apache/comet/shims/CometExprShim.scala
@YutaLin

YutaLin commented May 19, 2026

Copy link
Copy Markdown
Contributor Author

Hi @coderfender
Thanks for the review, i've made the change, please help me review again!

@YutaLin
YutaLin requested a review from coderfender May 19, 2026 04:30
@YutaLin

YutaLin commented May 21, 2026

Copy link
Copy Markdown
Contributor Author

Hi @coderfender
Could you help me check this again? Thanks!

@coderfender

Copy link
Copy Markdown
Contributor

Sorry for the delay. Let me look into it shortly @YutaLin

@andygrove

Copy link
Copy Markdown
Member

Thanks @YutaLin. This looks good. Could you fix the conflicts?

YutaLin added 3 commits June 2, 2026 18:23
…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
@YutaLin

YutaLin commented Jun 2, 2026

Copy link
Copy Markdown
Contributor Author

Hi @andygrove,
thanks for the review! I've fix the conflict!

@coderfender

Copy link
Copy Markdown
Contributor

Seems like there are still some more conflicts @YutaLin .

…expression_encode

# Conflicts:
#	docs/source/user-guide/latest/expressions.md
@YutaLin

YutaLin commented Jun 4, 2026

Copy link
Copy Markdown
Contributor Author

Hi @coderfender
Thanks for pointing that out. The previous conflict had already been resolved and reviewed. Since the PR wasn't merged immediately afterward, additional changes on the main branch introduced a new conflict. I've now merged the latest main branch and resolved that conflict as well.

Could you please take another look when you have a chance?

@coderfender

Copy link
Copy Markdown
Contributor

Sure. Yeah have been working on some doc changes causing merge conflicts. I will review it shortly

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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.getBytes

When 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.

@comphead

comphead commented Jul 2, 2026

Copy link
Copy Markdown
Contributor

is this PR related to apache/datafusion#21331 ?

YutaLin added 2 commits July 23, 2026 16:38
…expression_encode

# Conflicts:
#	spark/src/main/scala/org/apache/comet/serde/strings.scala
@YutaLin

YutaLin commented Jul 24, 2026

Copy link
Copy Markdown
Contributor Author

Hi @andygrove
Thanks for reviewing this. I kept the cast lowering for valid UTF-8, documented the malformed UTF-8 limitation next to the implementation and in the compatibility guide, and captured the reproduction under #4764. I also loosened the 4.x StaticInvoke match and fixed the trailing newline.

@andygrove

Copy link
Copy Markdown
Member

Note on this review: this was generated by an LLM (Claude Code) at my request while I worked through a review backlog. I have not verified the individual findings myself. Please treat everything below as suggestions to evaluate rather than as authoritative review feedback, and push back on anything that is wrong or already handled.

Lowering encode(str, 'utf-8') to a CAST(string AS binary) is a neat way to get this for free, and the CometExprShimCommon trait shared across the 4.x shims is the right structure for the StaticInvoke rewrite. The SQL fixture covers empty strings, NULL, multibyte, and the mixed-case charset literal, which is good.

Three things.

The malformed-UTF-8 divergence is not reported to users

Spark replaces malformed bytes during encode, while the cast lowering preserves them. That is a wrong answer, not a fallback, and right now the only record of it is an ignore(...) line in encode.sql plus a code comment.

encode should report Incompatible(Some(...)) for this, or at minimum the divergence needs to appear on the string compatibility page so a user can find it. As it stands, expressions.md will show encode as supported with no caveat and a user with dirty string data gets silently different bytes. That is the exact case Comet's Incompatible mechanism exists for.

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 Charset.forName accepts UTF8, utf8, and unicode-1-1-utf-8 as aliases for UTF-8, and Spark accepts whatever Charset.forName accepts. So encode(s, 'UTF8') falls back here even though it is exactly the case this PR handles.

Would Try(Charset.forName(name).name() == "UTF-8").getOrElse(false) be better? It handles aliases, and an invalid charset name naturally falls back to Spark, which then raises the proper INVALID_PARAMETER_VALUE.

Duplication between the 3.4 and 3.5 shims

The additions to spark-3.4/CometExprShim.scala and spark-3.5/CometExprShim.scala are identical. The 4.x side already shares through CometExprShimCommon. Is there a spark-3.x shared location that could host the same dispatch for 3.4 and 3.5? Two copies of a three-line match is minor now, but it is the pattern that makes per-version gaps show up later, and CI only lints a subset of the version matrix.

@andygrove andygrove added enhancement New feature or request area:expressions Expression evaluation labels Sep 6, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: encode fell 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 FF as EFBFBD, while the new lowering preserves FF. 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 its StaticInvoke replacement. 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.

@dwsmith1983

Copy link
Copy Markdown
Contributor

#3183 was closed by #5037. encode now goes through the codegen dispatcher on every profile (3.x via #5037, 4.x via the StaticInvoke catch-all from #5692), so it matches Spark on malformed UTF-8 too, which the cast lowering here doesn't. If the native cast is still worth it as a speedup, it needs a rebase and numbers against the dispatcher; otherwise I think this can close.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: At the specified base, encode falls 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 FF as EFBFBD, while the new cast preserves FF. 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-8 case-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 its StaticInvoke replacement. Both serialize the value through CometCast.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.

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

Labels

area:expressions Expression evaluation enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Feature] Support Spark expression: encode

6 participants