Skip to content

test: move expression-level codegen dispatcher tests to Comet SQL tests #6637

Description

@andygrove

Part of #6615.

What / Why

CometCodegenSuite tests the JVM codegen dispatcher, which compiles Spark's own doGenCode into an Arrow batch kernel. Most of it is dispatcher internals and stays in Scala: ScalaUDF registration, kernel caching, Nondeterministic state, TaskContext, memory lifecycle and explain text. But about a third of the suite tests how a built-in Spark expression routes through the dispatcher, or what it returns. expect_dispatch / expect_native / expect_fallback can express those now. Static analysis on main at 3bc2faa found the following:

Suite Covered Convert Blocked Split Keep Total
CometCodegenSuite 1 23 2 2 61 89

A built-in parent that always dispatches can stand in for an identity ScalaUDF. CometScalaUDF.emitJvmCodegenDispatch binds the whole subtree under the dispatched root: only column references become data arguments, and every inner node is reported as dispatched. Several serdes always dispatch:

  • map(...), map_concat
  • transform, transform_values, map_filter, zip_with, aggregate
  • hypot, nanvl, pmod, add_months, Invoke

So map('k', array_max(flatten(a))) puts array_max and flatten in one kernel over the nested input, the same as idBin(array_max(flatten(a))), and expect_dispatch(array_max, flatten) proves it. The same trick lets the CometCodegenHOFSuite tests in #6628 convert fully rather than split.

Work

Line numbers refer to CometCodegenSuite.

JSON routing

  • New: json/routing_json_{disabled,enabled,opt_in,opt_in_enabled}.sql, one file per allowIncompatible × codegen cell.
    • Covers L80 (json_array_length) and the two L151 instances (from_json, to_json). Issue Test dual-impl (native + codegen-dispatch) expressions consistently across the full routing matrix #4616's routing-matrix retrofit never reached JSON.
    • Use the Scala inputs: '[1,2,3]', '[]', 'not an array', NULL, and the struct/array JSON rows.
    • For from_json, use expect_fallback(from_json: spark.comet.exec.scalaUDF.codegen.enabled=false), expect_dispatch(from_json) or expect_native(from_json) per cell.
    • On Spark 4.x, json_array_length and to_json can't be named in expect_native / expect_dispatch. Both lower to StaticInvoke / Invoke there. Spark4xCometExprShim rebuilds the original expression but copies back only the fallback reasons, so the coverage tags are lost.
      • Use expect_fallback where a cell falls back, which works on 4.x, and a plain query elsewhere.
      • A fully Comet plan with the dispatcher off proves the native path. A dispatcher-off fallback with allowIncompatible=false proves there is no native path, and then a fully Comet plan with the dispatcher on proves dispatch.
      • Only "opt-in on, dispatcher on" for natively supported shapes stays unpinned on 4.x. A MaxSparkVersion: 3.5 copy could pin it on 3.x.
  • struct/json_to_structs.sql: change L24 and L28 from spark_answer_only to expect_dispatch(from_json).

Arrays

  • array/sequence.sql
    • L482: expect_native(sequence) on the explicit-step and default-step column queries.
    • L530: a d DATE table with sequence(d, DATE'2024-06-01', INTERVAL 1 MONTH) as expect_dispatch(sequence).
    • Upgrade L182-189 to expect_dispatch(sequence).
  • array/routing_arrays_{enabled,disabled,opt_in}.sql (L514): sequence(x, y, CASE WHEN x <= y THEN 2 ELSE -2 END).
  • array/transform.sql (L1815, L1832): transform(sequence(0, 15), i -> i * 10) and a named_struct variant, both expect_dispatch(transform) with no FROM, so 16 elements land in a one-row batch (JVM codegen dispatcher miscompiles map-typed (MapType) output #4539).
  • array/dispatch_nested_input.sql (new, spark.comet.batchSize=16): L1863, L1884, L1934, L1947, L1960, L1976.
    • map('k', array_distinct(s)) over ARRAY<STRUCT<a INT, b STRING>>.
    • map('k', array_max(flatten(a))) over BINARY, STRING, DECIMAL(10,2) and DECIMAL(30,2) nested arrays with inner NULLs, with expect_dispatch(array_max, flatten).
    • L1884's 100 deterministic rows are generated in SQL as one file.
  • array/dispatch_collection_output_sizes.sql (new, spark.comet.batchSize=8): L2324, 10 instances.
    • A 256-row single-file seed table, with one query per output type built by always-dispatch HOFs.

Strings, math, datetime, decimal, misc

  • string/upper.sql (L2023): expect_dispatch(upper, substring) on upper(substring(s, 1, 2)).
  • string/upper_ansi.sql (new, L2033): expect_error(DIVIDE_BY_ZERO) on IF(flag, upper(substring('abc', CAST(1L DIV 0L AS INT), n)), NULL), plus a sentinel of the same shape (Codegen dispatcher single-ordinal null short-circuit swallows ANSI errors from literal subtrees #5608).
  • string/lpad_binary.sql + _dispatch_disabled.sql (new, L2357): lpad(b, 8, unhex('FF')) over binary.
  • string/replace_dispatch_boundary.sql + _ansi.sql (new, L701)
    • Eight replace shapes as expect_dispatch(replace).
    • The ANSI null short-circuit shapes go in the ANSI file.
  • string/is_valid_utf8.sql (new, MinSparkVersion: 4.0, L2379): expect_dispatch(invoke) on is_valid_utf8(s) and make_valid_utf8(s), including a malformed value.
  • math/pmod.sql (L2066)
    • Add the (NULL, NULL) row, and make L26-27 expect_dispatch(pmod).
    • Add map('k', a + b) as expect_dispatch(add).
  • math/pmod_ansi.sql (new, L2087)
    • A (NULL, 0) sentinel row with expect_dispatch(pmod).
    • (7, 0) as expect_error(BY_ZERO), or use a 4.0 / 4.1 file pair to keep DIVIDE_BY_ZERO versus REMAINDER_BY_ZERO exact.
  • math/hypot.sql (L602): expect_dispatch(hypot, cast, checkoverflow, add) on hypot(CAST(amount + amount AS DOUBLE), 4.0D) over a DECIMAL(10,2) column.
  • math/hypot_dispatch_disabled.sql (new, L796): expect_fallback(hypot: spark.comet.exec.scalaUDF.codegen.enabled=false).
  • datetime/add_months_ansi.sql (new, L1992, Codegen dispatcher: whole-tree NullIntolerant short-circuit suppresses ANSI errors, plus TIME type gaps between canHandle and the runtime dispatcher #5218): expect_error(CAST_INVALID_INPUT) on add_months(CAST(s AS DATE), i), plus an expect_dispatch(add_months, cast) sentinel over valid rows.
  • decimal/decimal_ops.sql (L576): expect_native(checkoverflow, add) on a + b.
  • misc/dispatch_sliced_boolean.sql (new, L1066)
    • Set spark.sql.shuffle.partitions=1 and spark.comet.batchSize=100.
    • expect_dispatch(regexp_replace) on regexp_replace(IF(b, s, 'zz'), '1', 'y') over a grouped aggregate.

The conversions above assume the kernel accepts output types such as Map<String, Binary> and Map<String, Array<Struct>>. The UDF tests' map and struct outputs make that likely, but it hasn't been verified, so run the new fixtures before deleting the Scala tests.

Blocked or partly kept

  • L80 and the to_json instance of L151: blocked on the Spark 4.x coverage-tag loss described above. Pinning the mechanism on 4.x needs the shim to copy coverage tags back.
  • L701: the malformed CAST(X'FF' AS STRING) and 256 KiB repeat(...) literal cases need ConstantFolding (test: let Comet SQL tests keep ConstantFolding, exclude optimizer rules and change configs mid-file #6618), or the 262,144-character literal spelled out.
  • L2379: the 3.4/3.5 half (no SQL function reaches CometInvoke there) and the Invoke-with-arguments case on a folded literal target stay in Scala.

Stays in Scala

  • ScalaUDF registration and type surface: L978, L987, L1086-L1155, L1296 (10 instances), L1307-L1416, L1545-L1771.
  • UDF-specific semantics: L302, L498, L1030, L1075, L1440, L1461.
  • Kernel cache and Nondeterministic state: L818, L847, L864, L895, L921, L947, L1487, L1505.
  • Dispatcher counters: L317, L333.
  • Direct kernel and serde tests: L240, L2127 (TIME input, which has no SQL route), L2336, L2396.
  • Memory lifecycle: L270.
  • Explain, COMET-INFO and coverage-stat text: L357, L544, L626, L684.
  • Test-helper self-tests: L401, L430.
  • Dataset API and DSv2 catalog functions: L1040, L2441.

Already covered: L1802, by map/map_concat.sql L46-47.

Findings

  • The Codegen dispatcher: whole-tree NullIntolerant short-circuit suppresses ANSI errors, plus TIME type gaps between canHandle and the runtime dispatcher #5218 short-circuit tests L2066 and L2087 never generate the short-circuit. They wrap pmod / a + b in idInt(...), and a ScalaUDF root is not NullIntolerant. A bare pmod(a, b) does take the path, so the fixture versions are stronger.
  • L302's comment describes the wrong mechanism, and its claim that Concat(null, x) = x is wrong for Spark.
  • L847 never runs more than one batch per partition.
  • L864's premise doesn't hold: the table read forces columns nullable.
  • Explain on Spark 4.x misreports dispatched json_array_length and to_json. It likely does the same for parse_url, which goes through the same shim code. A dispatched call shows up as a native staticinvoke / invoke, and the COMET-INFO dispatcher line omits it. That is a product issue, not just a fixture gap. It was found by reading the code; worth confirming and filing separately.
  • Fixture headers claim a mechanism nothing asserts: sequence.sql, json_array_length_default.sql, hypot.sql, nanvl.sql, pmod.sql, map_concat.sql, create_map.sql, transform.sql.

Done when

  • Every test in the table below is accounted for. Either a fixture covers it (the same inputs, configs and assertion strength, or stronger), the fixture that already covers it is cited, or it stays in Scala for the reason given above.
  • Each converted Scala test is deleted in the PR that adds its coverage. A suite left empty is deleted and removed from the suite lists in .github/workflows/pr_build_linux.yml and pr_build_macos.yml.
  • The new and changed fixtures pass on every Spark profile. Apply run-all-spark-profiles, because pull requests run only the default profile.
Every Scala test in scope, with its verdict and target fixture
Suite Line Test Verdict Target
CometCodegenSuite 80 json_array_length routing follows native opt-in and dispatcher settings blocked: other json/routing_json_disabled.sql (new); json/routing_json_enabled.sql (new); json/routing_json_opt_in.sql (new); json/routing_json_opt_in_enabled.sql (new)
CometCodegenSuite 151 from_json routing follows native opt-in and dispatcher settings (loop at L133 instance 1 of 2; template "$name routing follows native opt-in and dispatcher settings") convert json/routing_json_disabled.sql (new); json/routing_json_enabled.sql (new); json/routing_json_opt_in.sql (new); json/routing_json_opt_in_enabled.sql (new); struct/json_to_structs.sql
CometCodegenSuite 151 to_json routing follows native opt-in and dispatcher settings (loop at L133 instance 2 of 2; template "$name routing follows native opt-in and dispatcher settings") blocked: other json/routing_json_disabled.sql (new); json/routing_json_enabled.sql (new); json/routing_json_opt_in.sql (new); json/routing_json_opt_in_enabled.sql (new); struct/structs_to_json.sql
CometCodegenSuite 240 codegen kernel round-trips CalendarIntervalType keep
CometCodegenSuite 270 a closed allocateOutput vector releases all of its memory keep
CometCodegenSuite 302 ScalaUDF over concat(c1, c2) suppresses the null short-circuit keep
CometCodegenSuite 317 disabled mode bypasses the dispatcher keep
CometCodegenSuite 333 schema exceeding spark.sql.codegen.maxFields falls back to Spark keep
CometCodegenSuite 357 explain.codegen.enabled surfaces routed expressions in COMET-INFO keep
CometCodegenSuite 401 checkSparkAnswerAndImpl pins the mechanism and fails when the claim is wrong keep
CometCodegenSuite 430 an expression nested inside a dispatched subtree is classified as dispatched keep
CometCodegenSuite 482 sequence with leaf integral args runs natively convert array/sequence.sql
CometCodegenSuite 498 sequence with zero-arg UDF stop routes through the dispatcher keep
CometCodegenSuite 514 sequence with non-leaf integral args routes through the dispatcher convert array/routing_arrays_enabled.sql; array/routing_arrays_disabled.sql; array/routing_arrays_opt_in.sql
CometCodegenSuite 530 sequence with date element type routes through the dispatcher convert array/sequence.sql
CometCodegenSuite 544 expression coverage stats split native from codegen-dispatch expressions keep
CometCodegenSuite 576 expression coverage stats survive the decimal promotion rewrite convert decimal/decimal_ops.sql
CometCodegenSuite 602 codegen dispatch coverage survives the decimal promotion rewrite convert math/hypot.sql
CometCodegenSuite 626 tags copied onto the shared TrueLiteral do not leak into unrelated plans keep
CometCodegenSuite 684 expression coverage stats count nothing when the plan falls back entirely keep
CometCodegenSuite 701 replace compatibility boundary cases stay on JVM codegen dispatcher split (rest: #6618) string/replace_dispatch_boundary.sql (new); string/replace_dispatch_boundary_ansi.sql (new)
CometCodegenSuite 796 codegen dispatch fallback reasons name the expression convert math/hypot_dispatch_disabled.sql (new)
CometCodegenSuite 818 dispatcher caches the compiled kernel across batches of one query keep
CometCodegenSuite 847 per-partition kernel preserves Nondeterministic state across batches keep
CometCodegenSuite 864 same UDF over nullable and non-nullable columns gets distinct kernels with independent state keep
CometCodegenSuite 895 Nondeterministic state persists across nullability flips within a partition keep
CometCodegenSuite 921 Nondeterministic state persists across two ScalaUDFs in one task keep
CometCodegenSuite 947 per-task cache isolates UDF state across sequential task runs in one session keep
CometCodegenSuite 978 registered string ScalaUDF routes through dispatcher keep
CometCodegenSuite 987 registered Java UDF1 routes through dispatcher keep
CometCodegenSuite 1030 boolean ScalaUDF argument sliced by an aggregate keeps its values keep
CometCodegenSuite 1040 typed Dataset.filter reads sliced booleans from a native aggregate keep
CometCodegenSuite 1066 dispatched regexp_replace reads a sliced boolean with its own values convert misc/dispatch_sliced_boolean.sql (new)
CometCodegenSuite 1075 struct ScalaUDF input with a sliced boolean child keeps its values keep
CometCodegenSuite 1086 multi-arg ScalaUDF over string + literal routes through dispatcher keep
CometCodegenSuite 1097 ScalaUDF as a child of a native Spark expression keep
CometCodegenSuite 1110 composed ScalaUDFs outer(inner(s)) fuse into one kernel keep
CometCodegenSuite 1126 ScalaUDFs of different types compose: isShort(len(s)) keep
CometCodegenSuite 1139 three-deep ScalaUDF composition lvl3(lvl2(lvl1(s))) keep
CometCodegenSuite 1155 multi-column ScalaUDF composition join(upperU(c1), lowerU(c2)) keep
CometCodegenSuite 1296 identity ScalaUDF on ${c.label} routes through dispatcher [instances: Boolean, Byte, Short, Int, Long, Float, Double, Date, Timestamp, TimestampNTZ] keep
CometCodegenSuite 1307 ScalaUDF returning a different type than its input keep
CometCodegenSuite 1319 ScalaUDF returning BinaryType keep
CometCodegenSuite 1331 ScalaUDF on BinaryType keep
CometCodegenSuite 1344 ScalaUDF returning ArrayType(StringType) keep
CometCodegenSuite 1361 ScalaUDF returning ArrayType(IntegerType) keep
CometCodegenSuite 1374 zero-column ScalaUDF produces one row per input row keep
CometCodegenSuite 1403 ScalaUDF over Decimal(18, 9) routes through the unscaled-long fast path keep
CometCodegenSuite 1416 ScalaUDF over Decimal(38, 10) routes through the BigDecimal slow path keep
CometCodegenSuite 1440 ScalaUDF sees TaskContext.partitionId() per partition keep
CometCodegenSuite 1461 ScalaUDF sees TaskContext from fully-native parquet plan keep
CometCodegenSuite 1487 Rand seeded per partition across a multi-partition table keep
CometCodegenSuite 1505 ScalaUDF composed with reused scalar subquery across projection and filter keep
CometCodegenSuite 1545 ScalaUDF taking Seq[String] reads element by element keep
CometCodegenSuite 1558 ScalaUDF taking Seq[String] iterating all elements keep
CometCodegenSuite 1571 ScalaUDF taking Seq[Int] reads primitive elements keep
CometCodegenSuite 1582 ScalaUDF taking Seq[BigDecimal] hits short-precision decimal fast path keep
CometCodegenSuite 1605 ScalaUDF composes with struct-field access reading Struct<String, Int>.age keep
CometCodegenSuite 1622 ScalaUDF taking full Struct<String, Int> value (case class arg) keep
CometCodegenSuite 1643 ScalaUDF returning Struct<String, Int> (case class output) keep
CometCodegenSuite 1652 ScalaUDF taking Map<String, Int> keep
CometCodegenSuite 1663 ScalaUDF round-trips Map<Int, Int> (primitive key and value) keep
CometCodegenSuite 1680 ScalaUDF returning Map<String, Int> keep
CometCodegenSuite 1693 ScalaUDF taking Map<String, Seq> exercises nested composition keep
CometCodegenSuite 1710 ScalaUDF round-trips Array<Array> (nested array input + output) keep
CometCodegenSuite 1729 ScalaUDF round-trips Struct<name, items: Array> keep
CometCodegenSuite 1749 ScalaUDF round-trips Map<String, Array> (nested value both sides) keep
CometCodegenSuite 1771 ScalaUDF round-trips Map<String, Struct<x: Int, y: String>> keep
CometCodegenSuite 1802 constant-folded map_concat output round-trips every key through the kernel (#4539) covered map/map_concat.sql
CometCodegenSuite 1815 constant-folded array output writes every element past the pre-sized child (#4539) convert array/transform.sql
CometCodegenSuite 1832 constant-folded Array<Struct<Int, String>> writes struct fields past the pre-sized child (#4539) convert array/transform.sql
CometCodegenSuite 1863 array_distinct on Array<Struct<Int, String>> retains element identity across hash set convert array/dispatch_nested_input.sql (new)
CometCodegenSuite 1884 array_max(flatten(arr)) on Array<Array> with mixed null inner arrays convert array/dispatch_nested_input.sql (new)
CometCodegenSuite 1934 array_max(flatten(arr)) on Array<Array> with null inner Binary returns null convert array/dispatch_nested_input.sql (new)
CometCodegenSuite 1947 array_max(flatten(arr)) on Array<Array> with null inner String returns null convert array/dispatch_nested_input.sql (new)
CometCodegenSuite 1960 array_max(flatten(arr)) on Array<Array<DECIMAL(10,2)>> with null inner Decimal (short-precision fast path) convert array/dispatch_nested_input.sql (new)
CometCodegenSuite 1976 array_max(flatten(arr)) on Array<Array<DECIMAL(30,2)>> with null inner Decimal (long-precision slow path) convert array/dispatch_nested_input.sql (new)
CometCodegenSuite 1992 multi-input NullIntolerant tree does not swallow an ANSI error (#5218) convert datetime/add_months_ansi.sql (new)
CometCodegenSuite 2023 single-input NullIntolerant tree still short-circuits nulls convert string/upper.sql
CometCodegenSuite 2033 single-input short-circuit does not swallow an ANSI error from a foldable subtree (#5608) convert string/upper_ansi.sql (new)
CometCodegenSuite 2066 multi-input leaf-only NullIntolerant tree short-circuits nulls correctly (#5218) convert math/pmod.sql
CometCodegenSuite 2087 leaf-only short-circuit preserves ANSI remainder-by-zero behaviour (#5218) convert math/pmod_ansi.sql (new)
CometCodegenSuite 2127 TIME input column routes through the dispatcher (#5218) keep
CometCodegenSuite 2324 dynamically-sized ${c.label} output round-trips through codegen dispatch (#4539) [instances: Array with null elements, Array, Array with null elements, Array, Array, Map<Int, Int>, Map<String, Int>, Array<Array>, Map<Int, Array>, Array<Struct<Int, String>>] convert array/dispatch_collection_output_sizes.sql (new)
CometCodegenSuite 2336 dispatch falls back cleanly when the bound tree cannot be closure-serialized (#5573) keep
CometCodegenSuite 2357 unrecognized StaticInvoke routes through the dispatcher instead of falling back (#5575) convert string/lpad_binary.sql (new); string/lpad_binary_dispatch_disabled.sql (new)
CometCodegenSuite 2379 Invoke routes through the codegen dispatcher (#5575) split (rest: #6618) string/is_valid_utf8.sql (new)
CometCodegenSuite 2396 the dispatcher declines calls into code other than Spark's (#6425) keep
CometCodegenSuite 2441 calls to a DataSource V2 function run in Spark (#6425) keep
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions