Repository navigation
feat: support Spark encode expression via codegen dispatch #5037
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
bf06051
39321c2
f21c07f
7837b97
a4ec13f
dd77303
f4cdc33
ac9e6f8
2e3d9bb
17c2d65
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| @@ -0,0 +1,28 @@ | ||||||||||||||||||||||||||||||||||||||
| /* | ||||||||||||||||||||||||||||||||||||||
| * Licensed to the Apache Software Foundation (ASF) under one | ||||||||||||||||||||||||||||||||||||||
| * or more contributor license agreements. See the NOTICE file | ||||||||||||||||||||||||||||||||||||||
| * distributed with this work for additional information | ||||||||||||||||||||||||||||||||||||||
| * regarding copyright ownership. The ASF licenses this file | ||||||||||||||||||||||||||||||||||||||
| * to you under the Apache License, Version 2.0 (the | ||||||||||||||||||||||||||||||||||||||
| * "License"); you may not use this file except in compliance | ||||||||||||||||||||||||||||||||||||||
| * with the License. You may obtain a copy of the License at | ||||||||||||||||||||||||||||||||||||||
| * | ||||||||||||||||||||||||||||||||||||||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||||||||||||||||||||||||||||||||||||||
| * | ||||||||||||||||||||||||||||||||||||||
| * Unless required by applicable law or agreed to in writing, | ||||||||||||||||||||||||||||||||||||||
| * software distributed under the License is distributed on an | ||||||||||||||||||||||||||||||||||||||
| * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY | ||||||||||||||||||||||||||||||||||||||
| * KIND, either express or implied. See the License for the | ||||||||||||||||||||||||||||||||||||||
| * specific language governing permissions and limitations | ||||||||||||||||||||||||||||||||||||||
| * under the License. | ||||||||||||||||||||||||||||||||||||||
| */ | ||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||
| package org.apache.comet.serde | ||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||
| import org.apache.spark.sql.catalyst.expressions.Encode | ||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||
| /** | ||||||||||||||||||||||||||||||||||||||
| * Spark 3.x `encode(str, charset)` runs through the codegen dispatcher so Spark's own encoder | ||||||||||||||||||||||||||||||||||||||
| * handles charset selection and malformed-input behavior. Dual of `CometStringDecode`. | ||||||||||||||||||||||||||||||||||||||
| */ | ||||||||||||||||||||||||||||||||||||||
| object CometEncode extends CometCodegenDispatch[Encode] | ||||||||||||||||||||||||||||||||||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Could you enable the existing
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Ungated in 17c2d65, and you were right that the gate was pointed at exactly the wrong versions — The gate was written when both cases only reached the dispatcher via Results, Spark 3.5.9 / JDK 17 /
|
||||||||||||||||||||||||||||||||||||||
| case | arm | best | per row |
|---|---|---|---|
encode(utf-8) |
dispatch off (Spark fallback) | 62 ms | 58.9 ns |
| codegen dispatch | 61 ms | 58.0 ns | |
| Spark (Comet disabled) | 78 ms | 74.4 ns | |
| dispatch off (repeat) | 59 ms | 56.0 ns | |
to_binary(utf-8) |
dispatch off (Spark fallback) | 64 ms | 60.8 ns |
| codegen dispatch | 59 ms | 56.3 ns | |
| Spark (Comet disabled) | 77 ms | 73.2 ns | |
| dispatch off (repeat) | 64 ms | 61.2 ns |
1,024 rows (sub-batch): encode 5 ms dispatch vs 5 ms fallback vs 9 ms Spark; to_binary 5 / 4 / 9.
Reading them honestly
For encode the repeated baseline is 59 ms against the first baseline's 62 ms, so this machine's noise floor for that table is ~3 ms — larger than the 1 ms between dispatch-on and dispatch-off. The right conclusion is that dispatch is performance-neutral there, not that it is 1 ms faster. to_binary has a 0 ms spread between its two baselines and a 5 ms gap (~8%), which is outside the noise, so that one is a real if modest win.
What both show is that routing through the dispatcher costs nothing relative to falling the whole projection back, and that either Comet arm beats Comet-disabled Spark by ~1.25x. That is the expected shape: both arms run Spark's own Encode implementation, so what is being traded is bridge overhead against losing the operator.
The run produced no WARNING lines, which matters more than the timings — checkPlans verifies per case that the dispatch arm is fully Comet native, that the dispatcher actually compiled or cache-hit a kernel, and that the dispatch-off arm is not fully native. So on 3.5 the two arms really are different plans and CometEncode is doing the work.
The to_time(fmt) case is still gated on 4.1, correctly, and the skip note now lists only that.
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,94 @@ | ||
| -- Licensed to the Apache Software Foundation (ASF) under one | ||
| -- or more contributor license agreements. See the NOTICE file | ||
| -- distributed with this work for additional information | ||
| -- regarding copyright ownership. The ASF licenses this file | ||
| -- to you under the Apache License, Version 2.0 (the | ||
| -- "License"); you may not use this file except in compliance | ||
| -- with the License. You may obtain a copy of the License at | ||
| -- | ||
| -- http://www.apache.org/licenses/LICENSE-2.0 | ||
| -- | ||
| -- Unless required by applicable law or agreed to in writing, | ||
| -- software distributed under the License is distributed on an | ||
| -- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY | ||
| -- KIND, either express or implied. See the License for the | ||
| -- specific language governing permissions and limitations | ||
| -- under the License. | ||
|
|
||
| -- Tests for the SQL `encode(str, charset)` function (StringType, StringType) -> BinaryType. | ||
| -- | ||
| -- `encode` runs through the codegen dispatcher (Spark's own doGenCode inside the Comet | ||
| -- pipeline) so behavior matches Spark exactly across all supported charsets and across the | ||
| -- Spark 4.0 `legacyCharsets` / `legacyErrorAction` modes. This is the dual of the `decode` | ||
| -- codegen-dispatch path (#4465). | ||
| -- | ||
| -- Each `query` block runs checkSparkAnswerAndOperator, which fails if the expression fell back | ||
| -- to Spark instead of executing natively, so these are non-vacuous. | ||
| -- Config: spark.comet.exec.scalaUDF.codegen.enabled=true | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Added all three branches in f4cdc33.
Each of the error fixtures leads with a sentinel |
||
|
|
||
| statement | ||
| CREATE TABLE test_encode(s string) USING parquet | ||
|
|
||
| statement | ||
| INSERT INTO test_encode VALUES ('hello'), ('world'), (''), ('café'), (NULL) | ||
|
|
||
| -- Charset form over multiple charsets | ||
|
|
||
| query | ||
| SELECT encode(s, 'utf-8') FROM test_encode | ||
|
|
||
| query | ||
| SELECT encode(s, 'UTF-8') FROM test_encode | ||
|
|
||
| query | ||
| SELECT encode(s, 'UTF-16') FROM test_encode | ||
|
|
||
| query | ||
| SELECT encode(s, 'UTF-16BE') FROM test_encode | ||
|
|
||
| -- UTF-16 emits a byte-order mark and UTF-16LE does not, so the two differ in output, and | ||
| -- UTF-16LE differs from UTF-16BE in byte order. | ||
|
|
||
| query | ||
| SELECT encode(s, 'UTF-16LE') FROM test_encode | ||
|
|
||
| query | ||
| SELECT encode(s, 'UTF-32') FROM test_encode | ||
|
|
||
| query | ||
| SELECT encode(s, 'ISO-8859-1') FROM test_encode | ||
|
Comment on lines
+37
to
+59
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The description says six charsets including UTF-16LE. This block plus line 62 covers utf-8, UTF-8, UTF-16, UTF-16BE, ISO-8859-1 and US-ASCII. UTF-16LE is absent and worth adding for real: UTF-16 emits a BOM and UTF-16LE does not, so the byte output genuinely differs. UTF-32 is the remaining entry in VALID_CHARSETS with no coverage.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Restored UTF-16LE and added UTF-32 in f4cdc33, with a comment recording why UTF-16LE is not redundant with UTF-16 (BOM) or UTF-16BE (byte order). All seven |
||
|
|
||
| -- US-ASCII: use ASCII-only input so this file stays version-agnostic. Unmappable input under | ||
| -- US-ASCII substitutes on Spark 3.4/3.5 and throws on Spark 4.0+, so it is covered by the paired | ||
| -- fixtures encode_unmappable.sql and encode_unmappable_strict.sql instead. | ||
|
|
||
| statement | ||
| CREATE TABLE test_encode_ascii(s string) USING parquet | ||
|
|
||
| statement | ||
| INSERT INTO test_encode_ascii VALUES ('hello'), ('world'), (''), (NULL) | ||
|
|
||
| query | ||
| SELECT encode(s, 'US-ASCII') FROM test_encode_ascii | ||
|
|
||
| -- Literal inputs including NULL and empty string | ||
|
|
||
| query | ||
| SELECT encode('hello', 'utf-8'), encode('', 'utf-8'), encode(CAST(NULL AS STRING), 'utf-8') | ||
|
|
||
| -- Non-literal charset column: the charset is an ordinary child expression, not required to be | ||
| -- foldable. The dispatcher handles it because it runs Spark's own code. | ||
|
|
||
| statement | ||
| CREATE TABLE test_encode_charset(s string, cs string) USING parquet | ||
|
|
||
| statement | ||
| INSERT INTO test_encode_charset VALUES ('hello', 'utf-8'), ('world', 'UTF-16BE'), ('café', 'ISO-8859-1'), (NULL, 'utf-8') | ||
|
|
||
| query | ||
| SELECT encode(s, cs) FROM test_encode_charset | ||
|
|
||
| -- Round-trip: decode(encode(x)) recovers the original string | ||
|
|
||
| query | ||
| SELECT decode(encode(s, 'UTF-8'), 'UTF-8') FROM test_encode | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,44 @@ | ||
| -- Licensed to the Apache Software Foundation (ASF) under one | ||
| -- or more contributor license agreements. See the NOTICE file | ||
| -- distributed with this work for additional information | ||
| -- regarding copyright ownership. The ASF licenses this file | ||
| -- to you under the Apache License, Version 2.0 (the | ||
| -- "License"); you may not use this file except in compliance | ||
| -- with the License. You may obtain a copy of the License at | ||
| -- | ||
| -- http://www.apache.org/licenses/LICENSE-2.0 | ||
| -- | ||
| -- Unless required by applicable law or agreed to in writing, | ||
| -- software distributed under the License is distributed on an | ||
| -- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY | ||
| -- KIND, either express or implied. See the License for the | ||
| -- specific language governing permissions and limitations | ||
| -- under the License. | ||
|
|
||
| -- encode() charset validation on Spark 3.4 and 3.5. | ||
| -- | ||
| -- On those versions `Encode` calls `String.getBytes(charset)` directly, so any charset the JVM | ||
| -- knows is accepted (including ones outside Spark 4.0's allowlist, e.g. windows-1252) and an | ||
| -- unknown name surfaces as `UnsupportedEncodingException` through `Platform.throwException`. | ||
| -- Spark 4.0+ instead routes through `CharsetProvider.forName`, which rejects anything outside | ||
| -- VALID_CHARSETS unless `spark.sql.legacy.javaCharsets` is set; those two branches are covered by | ||
| -- encode_invalid_charset_strict.sql and encode_legacy_charsets.sql. | ||
| -- MaxSparkVersion: 3.5 | ||
| -- Config: spark.comet.exec.scalaUDF.codegen.enabled=true | ||
|
|
||
| statement | ||
| CREATE TABLE test_encode_charset_legacy(s string) USING parquet | ||
|
|
||
| statement | ||
| INSERT INTO test_encode_charset_legacy VALUES ('hello'), ('café'), (''), (NULL) | ||
|
|
||
| -- windows-1252 is not in Spark 4.0's VALID_CHARSETS but is a JVM charset, so 3.x accepts it. | ||
| -- This is also the sentinel query: it fails if `encode` did not execute natively, so the | ||
| -- expect_error query below cannot pass vacuously through an operator-level fallback. | ||
|
|
||
| query | ||
| SELECT encode(s, 'windows-1252') FROM test_encode_charset_legacy | ||
|
|
||
| -- A charset name the JVM does not know at all. | ||
| query expect_error(UnsupportedEncodingException) | ||
| SELECT encode('hello', 'NO-SUCH-CHARSET') |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,50 @@ | ||
| -- Licensed to the Apache Software Foundation (ASF) under one | ||
| -- or more contributor license agreements. See the NOTICE file | ||
| -- distributed with this work for additional information | ||
| -- regarding copyright ownership. The ASF licenses this file | ||
| -- to you under the Apache License, Version 2.0 (the | ||
| -- "License"); you may not use this file except in compliance | ||
| -- with the License. You may obtain a copy of the License at | ||
| -- | ||
| -- http://www.apache.org/licenses/LICENSE-2.0 | ||
| -- | ||
| -- Unless required by applicable law or agreed to in writing, | ||
| -- software distributed under the License is distributed on an | ||
| -- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY | ||
| -- KIND, either express or implied. See the License for the | ||
| -- specific language governing permissions and limitations | ||
| -- under the License. | ||
|
|
||
| -- encode() charset validation in Spark 4.0's default mode. | ||
| -- | ||
| -- `CharsetProvider.forName` permits only us-ascii, iso-8859-1, utf-8, utf-16be, utf-16le, utf-16 | ||
| -- and utf-32 unless `spark.sql.legacy.javaCharsets` is enabled, and otherwise raises | ||
| -- `INVALID_PARAMETER_VALUE.CHARSET`. On Spark 3.4/3.5 a JVM charset outside that list is simply | ||
| -- accepted; see encode_invalid_charset.sql. Enabling the legacy flag restores that behavior on | ||
| -- 4.0+; see encode_legacy_charsets.sql. | ||
| -- MinSparkVersion: 4.0 | ||
| -- Config: spark.comet.exec.scalaUDF.codegen.enabled=true | ||
|
|
||
| statement | ||
| CREATE TABLE test_encode_charset_strict(s string) USING parquet | ||
|
|
||
| statement | ||
| INSERT INTO test_encode_charset_strict VALUES ('hello'), ('café'), (''), (NULL) | ||
|
|
||
| -- Sentinel: an allowlisted charset, asserting `encode` executes natively so the expect_error | ||
| -- queries below trip the kernel rather than being satisfied by an operator-level Spark fallback. | ||
|
|
||
| query | ||
| SELECT encode(s, 'utf-8') FROM test_encode_charset_strict | ||
|
|
||
| -- windows-1252 is a valid JVM charset but is outside VALID_CHARSETS. | ||
| query expect_error(INVALID_PARAMETER_VALUE) | ||
| SELECT encode('hello', 'windows-1252') | ||
|
|
||
| -- Column input takes the same path. | ||
| query expect_error(INVALID_PARAMETER_VALUE) | ||
| SELECT encode(s, 'windows-1252') FROM test_encode_charset_strict | ||
|
|
||
| -- A charset name the JVM does not know at all is rejected by the same check. | ||
| query expect_error(INVALID_PARAMETER_VALUE) | ||
| SELECT encode('hello', 'NO-SUCH-CHARSET') |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,45 @@ | ||
| -- Licensed to the Apache Software Foundation (ASF) under one | ||
| -- or more contributor license agreements. See the NOTICE file | ||
| -- distributed with this work for additional information | ||
| -- regarding copyright ownership. The ASF licenses this file | ||
| -- to you under the Apache License, Version 2.0 (the | ||
| -- "License"); you may not use this file except in compliance | ||
| -- with the License. You may obtain a copy of the License at | ||
| -- | ||
| -- http://www.apache.org/licenses/LICENSE-2.0 | ||
| -- | ||
| -- Unless required by applicable law or agreed to in writing, | ||
| -- software distributed under the License is distributed on an | ||
| -- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY | ||
| -- KIND, either express or implied. See the License for the | ||
| -- specific language governing permissions and limitations | ||
| -- under the License. | ||
|
|
||
| -- encode() under `spark.sql.legacy.javaCharsets=true` on Spark 4.0+. | ||
| -- | ||
| -- The flag is baked into the `Encode` node at analysis time and carried through to | ||
| -- `StaticInvoke(classOf[Encode], "encode", ...)` as a literal, so it reaches the codegen | ||
| -- dispatcher and `CharsetProvider.forName` accepts any JVM charset again. Without the flag the | ||
| -- same queries raise `INVALID_PARAMETER_VALUE.CHARSET`; see encode_invalid_charset_strict.sql. | ||
| -- MinSparkVersion: 4.0 | ||
| -- Config: spark.comet.exec.scalaUDF.codegen.enabled=true | ||
| -- Config: spark.sql.legacy.javaCharsets=true | ||
|
|
||
| statement | ||
| CREATE TABLE test_encode_legacy_charsets(s string) USING parquet | ||
|
|
||
| -- All rows are representable in windows-1252 ('é' is 0xE9), so the strict coding-error action | ||
| -- that stays in effect here does not fire. | ||
| statement | ||
| INSERT INTO test_encode_legacy_charsets VALUES ('hello'), ('café'), (''), (NULL) | ||
|
|
||
| query | ||
| SELECT encode(s, 'windows-1252') FROM test_encode_legacy_charsets | ||
|
|
||
| query | ||
| SELECT encode('café', 'windows-1252'), encode('hello', 'Shift_JIS') | ||
|
|
||
| -- The allowlisted charsets keep working with the flag on. | ||
|
|
||
| query | ||
| SELECT encode(s, 'UTF-16LE') FROM test_encode_legacy_charsets |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,50 @@ | ||
| -- Licensed to the Apache Software Foundation (ASF) under one | ||
| -- or more contributor license agreements. See the NOTICE file | ||
| -- distributed with this work for additional information | ||
| -- regarding copyright ownership. The ASF licenses this file | ||
| -- to you under the Apache License, Version 2.0 (the | ||
| -- "License"); you may not use this file except in compliance | ||
| -- with the License. You may obtain a copy of the License at | ||
| -- | ||
| -- http://www.apache.org/licenses/LICENSE-2.0 | ||
| -- | ||
| -- Unless required by applicable law or agreed to in writing, | ||
| -- software distributed under the License is distributed on an | ||
| -- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY | ||
| -- KIND, either express or implied. See the License for the | ||
| -- specific language governing permissions and limitations | ||
| -- under the License. | ||
|
|
||
| -- encode() over characters the target charset cannot represent, on Spark 3.4 and 3.5. | ||
| -- | ||
| -- On those versions `Encode` is a plain BinaryExpression whose eval is | ||
| -- `input.toString.getBytes(charset)`, and `String.getBytes` uses an encoder configured with | ||
| -- `CodingErrorAction.REPLACE`. Unmappable characters therefore become the charset's replacement | ||
| -- byte (0x3F, '?') rather than raising. Spark 4.0+ builds the encoder with | ||
| -- `CodingErrorAction.REPORT` and throws instead; that branch is covered by the paired fixture | ||
| -- encode_unmappable_strict.sql. | ||
| -- MaxSparkVersion: 3.5 | ||
| -- Config: spark.comet.exec.scalaUDF.codegen.enabled=true | ||
|
|
||
| statement | ||
| CREATE TABLE test_encode_unmappable(s string) USING parquet | ||
|
|
||
| statement | ||
| INSERT INTO test_encode_unmappable VALUES ('café'), ('中文'), ('hello'), (''), (NULL) | ||
|
|
||
| -- Column input: 'é' and the CJK characters are unmappable in US-ASCII and are substituted. | ||
|
|
||
| query | ||
| SELECT encode(s, 'US-ASCII') FROM test_encode_unmappable | ||
|
|
||
| -- 'é' is representable in ISO-8859-1 (0xE9) but the CJK characters are not. | ||
|
|
||
| query | ||
| SELECT encode(s, 'ISO-8859-1') FROM test_encode_unmappable | ||
|
|
||
| -- Literal input, pinning the substituted bytes rather than only asserting Spark/Comet agreement. | ||
|
|
||
| query | ||
| SELECT encode('café', 'US-ASCII') = CAST('caf?' AS BINARY), | ||
| encode('中文', 'US-ASCII') = CAST('??' AS BINARY), | ||
| encode('中文', 'ISO-8859-1') = CAST('??' AS BINARY) |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
The Implementation cell stays
—, which is the correct generated value on 4.x:encodeis registered in theFunctionRegistryon 3.4, 3.5 and 4.0 (FunctionRegistry.scala:523,:527,:546respectively), but on 4.x it is RuntimeReplaceable so no serde is keyed onclassOf[Encode]in that profile andGenerateDocsemits the placeholder. On 3.4 and 3.5 the same cell becomesCodegen dispatchbecauseCometEncodeis aCometCodegenDispatch(GenerateDocs.scala:194).decodealready has this shape, so the committed value is consistent. Just regenerate withdev/generate-release-docs.shrather than hand-editing the row.There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Verified by regeneration rather than by hand. I ran
GenerateDocs(the same invocationdev/generate-release-docs.shuses) against a scratch copy ofdocs/source/user-guide/latestunder three profiles:encodecellCodegen dispatch——The committed cell matches the 4.x generation exactly, so the file is unchanged. The 3.5 run is also a useful independent check that
CometEncodereally is keyed onclassOf[Encode]in that profile.One thing the run surfaced that is out of scope here: the tracked file has pre-existing drift on unrelated rows. Under 4.0 and 4.1 the generator emits
Nativefor<<,>>,>>>(committed—) and—forto_char,to_varchar(committedCodegen dispatch); 4.1 additionally differs onhour,minute,second,make_timestamp. Since the published site regenerates into a temp tree at build time, this only affects the in-repo copy. I left it alone rather than adding unrelated churn to this PR, but it may be worth a separate cleanup, or a CI check that the committed column matches generation.