conversion_funcs Expression Audits#

Audit notes for expressions in this category that have been audited. Absence of an entry means the expression has not been audited yet, not that it is unsupported. See the user guide Spark Expression Support for current support status.

cast#

  • Spark 3.4.3 (audited 2026-05-27): identical to 3.5.8 modulo Cast.canUpCast refactored to delegate to UpCastRule.canUpCast.

  • Spark 3.5.8 (audited 2026-05-27): baseline. Cast(child, dataType, timeZoneId, evalMode); eval modes are LEGACY, ANSI, TRY. The legacy Cast.canCast matrix and the Cast.canAnsiCast matrix decide acceptance per type pair. Comet routes via CometCast (spark/src/main/scala/org/apache/comet/expressions/CometCast.scala) using a per-source-type support matrix that returns Compatible, Incompatible(reason), or Unsupported(reason); literal children are short-circuited to Compatible() so CometLiteral validates them. The serialized Cast proto carries datatype, evalMode, timezone (default UTC), allowIncompat (from spark.comet.expression.Cast.allowIncompatible), and isSpark4Plus. The native side (native/spark-expr/src/conversion_funcs/cast.rs) implements explicit per-eval-mode branches for narrowing numeric casts that match Spark’s overflow exceptions, and falls through to DataFusion cast_with_options(safe = !ANSI) for the rest.

  • Spark 4.0.1 (audited 2026-05-27): VariantType added; StringType literals replaced with _: StringType to accommodate collated strings. (TimestampType, ByteType|ShortType|IntegerType) added to canAnsiCast. NullIntolerant -> nullIntolerant: Boolean refactor. New ToPrettyString.BinaryFormatter semantics for Binary -> String are replicated natively via spark_binary_formatter. Numeric-to-numeric matrix unchanged.

  • Spark 4.1.1 (audited 2026-05-27): TimeType added; many TimeType arms in canCast/canAnsiCast. Geospatial GeographyType / GeometryType types added with their own conversion rules. Numeric-to-numeric matrix unchanged.

  • Known divergences and gaps:

    • CAST(<binary> AS STRING) decodes the bytes with decode_utf8_spark_lossy, replacing invalid UTF-8 with U+FFFD exactly as the JVM’s new String(bytes, UTF_8) does (the same decoder as the native shuffle). It matches Spark’s rendered output, but decoding is not byte-preserving, so it diverges from Spark for operations that read the underlying bytes: round-trips such as CAST(CAST(x AS STRING) AS BINARY), and value identity, since distinct ill-formed byte sequences all decode to U+FFFD and can therefore compare equal in Comet where Spark compares the raw bytes (#4764).

    • Spark 4.0 collated StringType is not explicitly guarded; pattern equality is expected to keep collated-string casts falling back, but there is no test (#4489; umbrella #2190).

    • Spark 4.1 TimeType casts have no explicit Unsupported arm; they fall back implicitly but do not appear in the auto-generated compatibility doc (#4490).

    • CAST(<map> AS <map>) runs natively via cast_map_to_map; CometCast.isSupported recurses into the key and value casts, so the map arm inherits whichever support level the inner casts report.

    • spark.sql.legacy.castComplexTypesToString.enabled=true is Unsupported for any array/map/struct-to-string cast, so those casts fall back to Spark while the flag is on. Comet only implements the default ({}-wrapped, NULL-rendering) formatting. The flag is internal from Spark 4.0 onward and defaults to false.

    • CAST(<float|double> AS DECIMAL) is Compatible: it rounds the shortest decimal string form of the value (matching Spark’s Double.toString + BigDecimal.setScale(HALF_UP) path) rather than the binary value, taking a guarded binary fast path for values provably far from a rounding tie. NaN and infinity are null even in ANSI mode, matching Spark, which swallows the NumberFormatException from Decimal(double) before the ANSI overflow check. The support level carries a note: Double.toString only emits the shortest round-trip form on JDK 19 and later, so on older JDKs Spark’s own result can differ for a value whose shortest form lands exactly on a rounding tie at the target scale.

  • Spark registers the type-name conversion functions (bigint, binary, boolean, date, decimal, double, float, int, smallint, string, timestamp, tinyint) as cast aliases. Each lowers to the same Cast node, so Comet handles it via the cast implementation with the same compatibility profile.

  • Performance (tuned 2026-07-14, PR #4920): narrowing integer casts (spark_cast_int_to_int) map the values buffer in a single pass with Arrow unary/try_unary and carry the null buffer over untouched, replacing an element-by-element Option/Result iterator-collect. Up to 100x faster on narrowing casts. Benchmark: benches/cast_numeric.rs.

  • Performance (tuned 2026-07-15, PR #4940): float/double-to-decimal casts (cast_floating_point_to_decimal128) now convert in a single vectorized unary_opt pass that maps out-of-range values (NaN, infinity, precision overflow) to null, replacing the per-element Decimal128Builder loop. ANSI raises via an O(1) null-count check plus a rare element-wise rescan. 15-36% faster with no regression on any shape. Benchmark: benches/cast_float_to_decimal.rs.