apache / apache/datafusion-comet
[EPIC] Criterion bench coverage for all native expressions
- Dominant language
- Scala
- Stars
- 1.3k
- Forks
- 373
- Avg merge
- 2d 4h
- Merged PRs (30d)
- 198
Description
## Background
Follow-up to #5363. The Scala end-to-end microbenchmarks conflate expression cost with scan and columnar-to-row transfer, so a Rust-side Criterion bench is the right tool for actually attributing per-row time to an expression. We already have a bench harness under `native/spark-expr/benches/`, but coverage is ad hoc.
This issue is the catalog: what native expressions exist, which have Criterion benches, which don't, and a proposal for filling the gap deliberately (not one bench per one expression — there are hundreds of thin wrappers around identical shapes).
## Scope: what counts as a "native expression"
Two sources reach native code from Comet:
1. **Comet-owned kernels** in `native/spark-expr/src/`. These have their own Rust implementations that Comet is responsible for.
2. **DataFusion / datafusion-spark passthroughs** dispatched by function name via `CometScalarFunction(...)` in `spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala`. Comet's responsibility here is limited to argument massaging and correctness gating — the kernel itself is upstream.
For (2), a bench inside Comet's tree is redundant with upstream benches for the kernel itself, but *is* useful when Comet wraps the passthrough with pre/post processing (input coercion, timezone handling, NULL semantics) or when the compat gate (`allow_incompat`) selects a different path than the plain DF call. The catalog below treats those distinctly.
## Existing Criterion benches
40 bench files under `native/spark-expr/benches/`:
`aggregate`, `array_size`, `arrays_overlap`, `base64`, `bloom_filter_agg`, `cast_binary_to_string`, `cast_decimal_to_string`, `cast_float_to_decimal`, `cast_float_to_string`, `cast_from_boolean`, `cast_from_string`, `cast_int_to_decimal`, `cast_int_to_timestamp`, `cast_non_int_numeric_timestamp`, `cast_numeric`, `cast_string_to_date`, `cast_string_to_timestamp`, `ceil`, `check_overflow`, `checked_arithmetic`, `conditional`, `date_trunc`, `date_trunc_array_fmt`, `decimal_div`, `decimal_rescale`, `floor`, `get_json_object`, `make_decimal`, `map_sort`, `normalize_nan`, `padding`, `parse_url`, `regexp_extract`, `round`, `to_csv`, `to_json`, `to_time`, `unhex`, `unscaled_value`, `wide_decimal`
Plus `native/core/benches/array_element_append.rs` and `perf.rs`.
## Coverage matrix — Comet-owned kernels
Categories mirror the directory layout under `native/spark-expr/src/`. `[x]` = bench exists, `[ ]` = gap. Passthrough entries are shown separately at the end.
### `agg_funcs`
- [x] `avg_decimal` — covered by `aggregate.rs`
- [x] `sum_decimal` — covered by `aggregate.rs`
- [ ] `sum_int`
- [ ] `avg` (non-decimal)
- [ ] `approx_percentile`
- [ ] `percentile`
- [ ] `correlation`
- [ ] `covariance` (pop, samp)
- [ ] `variance` (pop, samp) / `stddev` (pop, samp) — the Welford path
- [ ] `hll_plus_plus`
- [x] `bloom_filter_agg`
- [ ] `first`, `last`
- [ ] `collect_list`, `collect_set`
- [ ] `bit_and_agg`, `bit_or_agg`, `bit_xor_agg`
### `array_funcs`
- [x] `arrays_overlap`
- [x] `array_size` (`size`)
- [ ] `array_insert`
- [ ] `array_position`
- [ ] `array_slice`
- [ ] `arrays_zip`
- [ ] `flatten`
- [ ] `get_array_struct_fields`
- [ ] `list_extract`
- [ ] `sort_array`
- [ ] `array_join`, `array_max`, `array_min`, `array_remove`, `array_intersect`, `array_union`, `array_except`, `array_contains`, `array_transform`, `array_exists`, `array_forall`, `array_aggregate`, `array_sort`, `zip_with`, `sequence`, `shuffle`, `create_array`, `get_array_item`, `element_at` (all reachable via serde entries in `arrayExpressions`)
### `conditional_funcs`
- [x] `if_expr` — covered by `conditional.rs`
- [x] `case_when` — covered by `conditional.rs`
### `conversion_funcs` (cast)
Well-covered already:
- [x] `cast_numeric`, `cast_from_string`, `cast_from_boolean`, `cast_int_to_decimal`, `cast_float_to_decimal`, `cast_int_to_timestamp`, `cast_string_to_date`, `cast_string_to_timestamp`, `cast_non_int_numeric_timestamp`, `cast_float_to_string`, `cast_binary_to_string`, `cast_decimal_to_string`
Gaps:
- [ ] `cast_string_to_numeric` (there is `cast_from_string` but the string-to-numeric path should be exercised across all target widths)
- [ ] `cast_timestamp_to_string`, `cast_date_to_string`
- [ ] `cast_timestamp_to_date`, `cast_date_to_timestamp`
- [ ] `cast_decimal_to_decimal` (rescale is covered, but full cast between decimal widths is not)
- [ ] `trim` variants that route through `conversion_funcs/trim.rs` (leading/trailing/both, custom trim strings) — passthrough today, but Comet's serde does the coercion
### `csv_funcs`
- [x] `to_csv`
- [ ] `csv_to_structs` (from_csv) — the more expensive direction, unbenched
### `datetime_funcs`
- [x] `date_trunc`, `date_trunc_array_fmt`, `to_time`
- [ ] `timestamp_trunc`
- [ ] `unix_timestamp` (timestamp input, date input)
- [ ] `from_unix_time`
- [ ] `date_add`, `date_sub`, `date_diff`
- [ ] `date_from_unix_date`, `unix_date`
- [ ] `make_date`, `make_time`, `make_interval`, `make_timestamp`, `make_ym_interval`, `make_dt_interval`, `multiply_dt_interval`
- [ ] `hours` / `hour` / `minute` / `second` / `day_of_month` / `day_of_week` / `day_of_year` / `week_of_year` / `week_day` / `quarter` / `year` / `month` — most are thin, but `extract_date_part.rs` centralizes them and is worth one parameterized bench across the field list
- [ ] `next_day`, `last_day`, `add_months`, `months_between`
- [ ] `from_utc_timestamp`, `to_utc_timestamp`, `convert_timezone` — timezone handling has been a repeat perf hotspot
- [ ] `seconds_to_timestamp`, `micros_to_timestamp`, `millis_to_timestamp`
- [ ] `timestamp_add`, `timestamp_diff`
### `hash_funcs`
- [ ] `murmur3` — used per row on every hash-partitioned shuffle key; regressions here hit shuffle throughput directly
- [ ] `xxhash64`
### `json_funcs`
- [x] `to_json`, `get_json_object`
- [ ] `from_json`
- [ ] `json_array_length`
- [ ] `json_object_keys`, `schema_of_json`
### `map_funcs`
- [x] `map_sort`
- [ ] `map_extract` (`get_map_value`), `map_keys`, `map_values`, `map_entries`, `map_from_arrays`, `map_from_entries`, `map_concat`, `str_to_map`, `map_filter`, `transform_keys`, `transform_values`, `map_zip_with`, `create_map`
### `math_funcs`
- [x] `ceil`, `floor`, `round`, `unhex`, `checked_arithmetic`, `decimal_div`, `decimal_rescale`, `unscaled_value`, `wide_decimal`, `check_overflow`, `normalize_nan`, `make_decimal`
- [ ] `abs`
- [ ] `log`, `log10`, `log2`, `log1p`, `logarithm`
- [ ] `pow`, `hypot`
- [ ] `modulo` (`pmod`, `remainder`)
- [ ] `negative` / `unary_minus`
- [ ] `width_bucket`
- [ ] `conv`, `hex`
- [ ] `bround`
- [ ] `nanvl`
### `nondetermenistic_funcs`
Timing is meaningful for the RNG path even though the values are random:
- [ ] `rand`, `randn`, `rand_str`
- [ ] `uuid`
- [ ] `bernoulli_cell_sampler`, `shuffle`
- [ ] `monotonically_increasing_id`
### `predicate_funcs`
- [ ] `is_nan`
- [ ] `rlike` — critical because it has native Rust and JVM fallback modes; regression here shows up in every LIKE-heavy query
### `string_funcs`
- [x] `base64`, `regexp_extract`, `padding` (lpad/rpad)
- [ ] `unbase64` — has its own kernel
- [ ] `contains`, `starts_with`, `ends_with`
- [ ] `length`, `octet_length`, `bit_length` (thin, but common)
- [ ] `levenshtein`
- [ ] `regexp_extract_all`, `regexp_in_str`, `regexp_replace`
- [ ] `split`
- [ ] `upper`, `lower`, `init_cap` — these are `allow_incompat` gated; the two paths (built-in ICU vs Comet's) should both be benched
- [ ] `substring`, `substring_index`, `left`, `right`
- [ ] `overlay`
- [ ] `mask`, `sound_ex`, `format_number`, `format_string`
- [ ] `find_in_set`, `string_locate`
- [ ] `elt`
- [ ] `to_number`, `try_to_number`, `to_character`
- [ ] `reverse` (string variant; array variant separately)
- [ ] `empty2null`
### `struct_funcs`
- [ ] `create_named_struct`
- [ ] `get_struct_field`
### `url_funcs`
- [x] `parse_url`
### Misc / `static_invoke` / kernels
- [ ] `xpath` family (`XPathBoolean`, `XPathShort`, `XPathInt`, `XPathLong`, `XPathFloat`, `XPathDouble`, `XPathString`, `XPathList`) — 8 entries, likely one parameterized bench
- [ ] `bloom_filter_might_contain` (probe side; agg side is covered)
- [ ] `char_varchar_utils` (static_invoke) — pad/trim for CHAR/VARCHAR semantics
## Coverage matrix — DataFusion passthroughs
40 expressions dispatched by name via `CometScalarFunction`. The kernel itself is upstream. Comet-owned coverage here should focus on the passthroughs whose serde does non-trivial work or where Comet has a compat gate. Everything else can be left to upstream benches.
Passthroughs listed for completeness (from `QueryPlanSerde.scala`):
- **Math (trig / transcendentals):** `acos`, `acosh`, `asin`, `asinh`, `atan`, `atanh`, `cbrt`, `cos`, `cosh`, `cot`, `csc`, `degrees`, `exp`, `expm1`, `factorial`, `greatest`, `least`, `pi`, `radians`, `rint`, `sec`, `signum`, `sin`, `sinh`, `sqrt`, `tan`, `tanh`, `bin`. **Skip individually.** Add one parameterized "trig fastpath" bench that runs the whole set on Float64 with and without nulls, so any regression across the whole family surfaces at once.
- **Hash:** `crc32`, `md5`. **Bench both.** Small kernels but on the hot path for many workloads.
- **Strings:** `ascii`, `char`, `instr`, `space`, `trim`, `ltrim`, `rtrim`. **Bench `trim` variants** (Comet's serde does whitespace / custom-trim-char routing); others can go in the parameterized string-fastpath bench.
- **Arrays:** `array_distinct`, `array_repeat`. **Bench both** — non-trivial per-row work.
- **Bitwise:** `shiftrightunsigned`. Skip; trivial.
## Proposal
**Not** one bench per expression. Concretely:
1. **A shared harness** — `native/spark-expr/benches/common.rs` (file already exists but is empty) with helpers for building `RecordBatch`es with deterministic seeded data at a standard set of row counts (e.g. 8k / 64k / 512k) and typical null ratios (0%, 10%, 100%). Every new bench uses those builders. Fixes the reproducibility problem the Scala benches have.
2. **Parameterized "family" benches** for expressions that share a shape. Trig fastpath is one. Extract-date-part fields is another (14 date-part accessors → one bench with `BenchmarkId` per field). String-fastpath (`ascii`, `chr`, `space`, `length`, `octet_length`, `bit_length`) is another. Cuts total bench code by ~5×.
3. **Individual benches** for the expressions listed with `[ ]` above where perf is non-trivial, has a compat gate, or has been a historical hotspot. Prioritize in this order:
- **P0 (hot path for typical Spark workloads):** `murmur3`, `xxhash64`, `abs`, `upper`/`lower`/`init_cap` (both compat modes), `contains`/`starts_with`/`ends_with`, `regexp_replace`, `substring`, timezone conversions (`from_utc_timestamp`, `to_utc_timestamp`, `convert_timezone`), `unix_timestamp`, `timestamp_trunc`, `is_nan`, `rlike` (both native and JVM modes), the remaining aggregates in `agg_funcs`.
- **P1 (correctness-gated / has known variants):** `from_json`, `csv_to_structs`, cast timestamp/date crossings, decimal-to-decimal cast, `unbase64`, `mask`, `to_number` / `try_to_number`.
- **P2 (thin wrappers):** everything else — most `array_funcs`, `map_funcs`, `struct_funcs`, `xpath` (one parameterized bench for the family), `format_number`/`format_string`, `sound_ex`, `find_in_set`, `string_locate`, `elt`, `reverse`, `empty2null`, misc math.
4. **A coverage script** in CI — `native/spark-expr/benches/coverage.sh` (or a small Rust test) that reads the list of expression classes registered in `QueryPlanSerde.scala` and asserts each one either has a bench file, appears in a parameterized-family allowlist, or is on an explicit skip-list with a reason. This is what prevents the gap from reopening. Failing this gate should be a merge blocker.
5. **A README section** in `native/spark-expr/benches/` documenting the convention (harness, naming, parameterization rules, skip-list) so contributors adding a new expression know a bench is expected.
Rough size estimate: ~40 new bench files needed to cover the P0+P1 list, if each is ~150 lines and reuses the shared harness. That is one focused sprint of work, or drip-fed alongside the individual expression changes as they come up. The coverage script is what makes the drip-feed viable.
Explicitly out of scope: the Scala end-to-end microbenchmarks. Those have separate problems (see #5363) and Criterion benches don't replace them because they don't exercise the JNI or serde boundary.
Contributor guide
Assessment
This issue has not been assessed yet.