String Expressions#
BitLength#
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
BinaryTypeinput is not supported
Concat#
By default, Concat is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.expression.Concat.allowIncompatible=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
concat does not support non-UTF8_BINARY collations (https://github.com/apache/datafusion-comet/issues/2190)
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
CONCAT supports only string input parameters
ConcatWs#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
all arguments are foldable
Contains#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
Non-UTF8_BINARY collated operands are routed through the JVM codegen dispatcher (Spark’s own
doGenCode) because native comparison is byte-wise.
EndsWith#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
Non-UTF8_BINARY collated operands are routed through the JVM codegen dispatcher (Spark’s own
doGenCode) because native comparison is byte-wise.
GetJsonObject#
By default, GetJsonObject is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.expression.GetJsonObject.allowIncompatible=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
Spark allows single-quoted JSON and unescaped control characters which Comet does not support
Very long numbers near Jackson’s 1000-digit limit can depend on Spark’s recycled parser buffer, which Comet cannot reproduce from the input alone
Selected JSON integers outside the 64-bit range and very long numbers can lose precision or fail during Comet’s native JSON materialization
Some selected floating-point values can differ from Spark at decimal parsing or Java-version formatting boundaries
When a returned object or array contains duplicate keys, Spark preserves them while Comet’s native JSON materialization keeps only the last value
InitCap#
By default, InitCap is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.expression.InitCap.allowIncompatible=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
Treats hyphen as a word separator (e.g.
robert rose-smithproducesRobert Rose-Smithinstead of Spark’sRobert Rose-smith) (https://github.com/apache/datafusion-comet/issues/1052)
Levenshtein#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
Non-default (non-UTF8_BINARY) collated input. The native kernel compares raw bytes, so collation-aware comparison has no native path.
Like#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
LIKE with a custom escape character (only
\is supported natively)Non-UTF8_BINARY collated operands are routed through the JVM codegen dispatcher (Spark’s own
doGenCode) because native comparison is byte-wise.
Lower#
By default, Lower is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.caseConversion.enabled=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
Results can vary depending on locale and character set
OctetLength#
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
BinaryTypeinput is not supported
RLike#
The following cases use Comet’s native implementation by default:
A
UTF8_BINARYliteral pattern admitted by the plan-time compatibility analyzer is evaluated natively by default.
For applicable cases that are not selected for native execution automatically, RLike is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline) by default. Set spark.comet.expression.RLike.allowIncompatible=true to explicitly select Comet’s native implementation, which has the following differences from Spark:
For applicable literal patterns outside the automatically admitted subset, the native Rust regex engine may behave differently from Java regex.
RandStr#
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
lengthmust be a literal valuelengthmust be non-negativeseedmust be a literal value
RegExpReplace#
By default, RegExpReplace is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.expression.RegExpReplace.allowIncompatible=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
Regexp pattern may not be compatible with Spark
Reverse#
By default, Reverse is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.expression.Reverse.allowIncompatible=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
native reverse does not support arrays whose element type contains binary, struct, or map
reverse does not support non-UTF8_BINARY collations (https://github.com/apache/datafusion-comet/issues/2190)
StartsWith#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
Non-UTF8_BINARY collated operands are routed through the JVM codegen dispatcher (Spark’s own
doGenCode) because native comparison is byte-wise.
StringInstr#
For applicable cases that are not selected for native execution automatically, StringInstr is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline) by default. Set spark.comet.expression.StringInstr.allowIncompatible=true to explicitly select Comet’s native implementation, which has the following differences from Spark:
Spark evaluates non-UTF8_BINARY collated string input under its collation, while Comet’s native instr compares raw bytes
StringLPad#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
Scalar values are not supported for the
strargument.Only scalar values are supported for the
padargument.
StringRPad#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
Scalar values are not supported for the
strargument.Only scalar values are supported for the
padargument.
StringRepeat#
The following differences from Spark are always present and do not require any additional configuration:
A negative argument for the number of times to repeat throws an exception instead of returning an empty string as Spark does
StringReplace#
By default, StringReplace is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.expression.StringReplace.allowIncompatible=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
Produces different results from Spark when the search string is empty
StringSplit#
By default, StringSplit is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.expression.StringSplit.allowIncompatible=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
Regex engine differences between Java and Rust
StringTranslate#
By default, StringTranslate is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.expression.StringTranslate.allowIncompatible=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
DataFusion’s translate iterates over Unicode graphemes (Spark uses code points) and substitutes U+0000 instead of treating it as a deletion sentinel
StringTrim#
For applicable cases that are not selected for native execution automatically, StringTrim is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline) by default. Set spark.comet.expression.StringTrim.allowIncompatible=true to explicitly select Comet’s native implementation, which has the following differences from Spark:
Spark evaluates non-UTF8_BINARY collated string input under its collation, while Comet’s native trim with a trim string compares raw bytes
StringTrimLeft#
For applicable cases that are not selected for native execution automatically, StringTrimLeft is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline) by default. Set spark.comet.expression.StringTrimLeft.allowIncompatible=true to explicitly select Comet’s native implementation, which has the following differences from Spark:
Spark evaluates non-UTF8_BINARY collated string input under its collation, while Comet’s native ltrim with a trim string compares raw bytes
StringTrimRight#
For applicable cases that are not selected for native execution automatically, StringTrimRight is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline) by default. Set spark.comet.expression.StringTrimRight.allowIncompatible=true to explicitly select Comet’s native implementation, which has the following differences from Spark:
Spark evaluates non-UTF8_BINARY collated string input under its collation, while Comet’s native rtrim with a trim string compares raw bytes
SubstringIndex#
For applicable cases that are not selected for native execution automatically, SubstringIndex is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline) by default. Set spark.comet.expression.SubstringIndex.allowIncompatible=true to explicitly select Comet’s native implementation, which has the following differences from Spark:
Spark evaluates non-UTF8_BINARY collated string input under its collation, while Comet’s native substring_index compares raw bytes
UnBase64#
The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):
unbase64 with failOnError = true uses stricter RFC 4648 validation that is not yet implemented natively
unbase64 with a non-trivial child expression uses the JVM codegen dispatcher to preserve Spark’s short-circuit evaluation (native path is limited to column and literal children)
Upper#
By default, Upper is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set spark.comet.caseConversion.enabled=true to opt into Comet’s native implementation instead, which has the following differences from Spark:
Results can vary depending on locale and character set