Compatibility Guide#

Comet aims to provide consistent results with the version of Apache Spark that is being used.

This guide documents areas where Comet’s behavior is known to differ from Spark. Topics are grouped by subsystem:

  • Parquet: limitations when reading Parquet files.

  • Floating-point comparison: NaN and signed-zero handling in comparisons.

  • Regular expressions: differences between the Rust regexp crate and Java’s regex engine.

  • Operators: operator-level compatibility notes, including window functions and round-robin partitioning.

  • Expressions: per-expression compatibility notes, including cast.

  • JSON: choosing between the native and Spark-compatible engines for JSON expressions.

  • Spark versions: version-specific known issues and limitations.

Compatible by default, opt in to native#

Comet runs a Spark-compatible implementation of every supported expression by default. Some expressions also have a native implementation that can differ from Spark for certain inputs. These are not used unless you opt in by setting the relevant spark.comet.expression.<Name>.allowIncompatible=true config (a few use a dedicated config, noted per expression below), after which you accept the documented differences.

You can discover where a native opt-in is available for a specific query in the verbose extended explain output. A [COMET-INFO: ...] segment points at an available native path and does not mean the operator falls back to Spark. This is distinct from [COMET: ...], which records a reason an operator did fall back.

Native and codegen-dispatch implementations#

Some Spark expressions have two implementations in Comet:

  • A codegen-dispatch implementation that runs Spark’s own generated code for the expression inside Comet’s native pipeline (via the Arrow-direct codegen dispatcher). This produces byte-exact Spark results at the cost of one JNI round-trip per batch. It is gated globally by spark.comet.exec.scalaUDF.codegen.enabled (enabled by default); when the dispatcher is disabled, these expressions fall back to Spark.

  • A native (Rust / DataFusion) implementation that avoids the JNI round-trip but has known semantic differences from Spark for some inputs or patterns.

Because the codegen-dispatch path matches Spark exactly, Comet uses it by default. The native path is opt-in per expression via that expression’s spark.comet.expression.<ExprClassName>.allowIncompatible=true flag, which declares that you accept its differences from Spark. There is no global opt-in. When the native path is enabled but a specific input or pattern has no native implementation, Comet routes that case back through the codegen dispatcher rather than running something incompatible.

This is the model behind the regular expression and JSON families, which document their per-expression configs and the specific differences to expect.

This is distinct from expressions that have no codegen-dispatch path: there, the incompatible cases fall back to Spark by default, and allowIncompatible=true runs the native (incompatible) path instead. Aggregate functions such as mode are the main example, because the codegen dispatcher covers only scalar expressions; see the expression reference for which expressions have incompatible cases.

Strings with non-UTF-8 bytes#

Spark’s StringType can hold arbitrary bytes, including sequences that are not valid UTF-8 (for example CAST(X'FF' AS STRING)). Arrow’s string type requires valid UTF-8, so Comet cannot store the raw bytes natively. When Comet produces a string from arbitrary bytes (such as CAST(binary AS string) or a columnar shuffle), it decodes them the same way the JVM does (new String(bytes, UTF_8)), replacing each ill-formed sequence with the Unicode replacement character U+FFFD. Spark itself applies the identical replacement whenever such a string is materialized (collected, printed, or passed to most string functions), so the rendered result matches Spark.

Decoding is not byte-preserving, so results can differ from Spark for any operation that works on the underlying bytes rather than on the rendered text:

  • Round-trips. Spark keeps the original bytes, so CAST(CAST(X'FF' AS STRING) AS BINARY) returns X'FF', whereas Comet returns the UTF-8 encoding of U+FFFD (X'EFBFBD'). octet_length and hashing of such a string differ for the same reason.

  • Value identity. Decoding maps every ill-formed sequence onto the same U+FFFD, so two Spark strings that hold different bytes can become equal in Comet. For example, with b = X'FF', CAST(b AS STRING) = CAST(X'EFBFBD' AS STRING) is false in Spark (UTF8String compares the raw bytes) but true in Comet. Equality, joins, grouping, ordering, and byte-based string functions such as contains can therefore disagree with Spark when non-UTF-8 bytes are involved.

Both differences require string data that is not valid UTF-8, which does not occur for text read from Parquet or produced by string expressions. Consistent handling of invalid UTF-8 across all native string paths is tracked by #4764.

Separately, Comet’s native Parquet scan currently rejects string columns whose stored bytes are not valid UTF-8 rather than reading them like Spark (#4121).

ANSI-mode error classes and messages#

Under spark.sql.ansi.enabled=true, several native error paths raise the error at the correct input but with a different exception type, error class, SQLSTATE, or message text than Spark. Code that catches SparkException and only asserts on message substrings is unaffected; code that inspects the exception class, getCondition(), or the parameterised error class will observe divergence:

  • Byte / Short Add, Subtract, and Multiply overflow raises ARITHMETIC_OVERFLOW (for example byte overflow) where Spark raises BINARY_ARITHMETIC_OVERFLOW, and integral ARITHMETIC_OVERFLOW messages omit Spark’s try_ suggestion (#6217).

  • Wide-decimal arithmetic overflow, decimal divide-by-zero, and decimal-to-decimal cast overflow raise raw Arrow errors that bypass SparkErrorConverter and surface as CometNativeException rather than SparkArithmeticException with the proper error class and query context (#5072).

  • Spark 4.2 introduced additional ANSI arithmetic overflow behavior differences that Comet does not yet track (#4967).

Known result-value divergences#

The following native paths silently return values that differ from Spark for edge-case inputs. Most also have entries in the per-category expression pages linked above; they are collected here so users hunting an unexpected value have a single place to check:

  • CAST(string AS timestamp) and CAST(string AS timestamp_ntz) trim Unicode whitespace. Spark trims only the bytes 0x00-0x20 and 0x7F, so a value padded with an ASCII control byte parses in Spark and returns NULL in Comet, while a value padded with non-ASCII whitespace such as U+3000 returns NULL in Spark and parses in Comet (#5149).

  • Native RANGE window frames with an explicit PRECEDING / FOLLOWING offset diverge from Spark when the boundary arithmetic overflows for DATE or DECIMAL ORDER BY columns (#5022).

  • Ungrouped decimal SUM keeps an unbounded intermediate and checks the result precision only when a partial is written out or the sum is evaluated, which matches Spark’s whole-stage codegen path. Without codegen Spark buffers the aggregate in an UnsafeRow and latches as soon as a running sum leaves the precision. Comet falls back at precision 38 when codegen is disabled by spark.sql.codegen.wholeStage, by spark.sql.codegen.factoryMode=NO_CODEGEN on Spark 3.5+, by an imperative sibling aggregate, by an expression Spark cannot compile such as a lambda function, or by the spark.sql.codegen.maxFields limit on the aggregate’s own output and inputs, but not when Spark abandons codegen at runtime, after a compile failure under spark.sql.codegen.fallback or when the generated code exceeds spark.sql.codegen.hugeMethodLimit, where an intermediate overflow that later cancels out returns NULL (or raises under ANSI) in Spark but the recovered value in Comet.