apache / apache/datafusion-comet

[EPIC] Criterion bench coverage for all native expressions

Open
#5,396 1 comment 0 reactions 1 assignee Claimed by @coderfender View on GitHub
area:ci area:expressions enhancement EPIC performance test
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

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.