andygrove opened a new issue, #5363:
URL: https://github.com/apache/datafusion-comet/issues/5363
## Background
We now have ~25 Scala end-to-end microbenchmark suites for expressions,
almost all of which funnel through `runExpressionBenchmark` in
`spark/src/test/scala/org/apache/spark/sql/benchmark/CometBenchmarkBase.scala`.
I reviewed them to check whether they actually isolate expression evaluation
cost. They largely do not — a large and non-uniform share of each timed region
is Parquet scan plus columnar-to-row conversion of the result set, and there
are several fairness bugs that make individual rows misleading.
These benchmarks are what we use to justify optimization work and to catch
regressions, so it is worth reimplementing the harness rather than patching
individual suites.
## Problem 1: scan and result transfer are inside the timed region
Every timed case is:
```scala
spark.sql("SELECT expr(c1) FROM parquetV1Table").noop()
```
Two costs that are not expression evaluation are inside the timer:
**Parquet scan.** Read on every iteration, and critically the two arms use
*different* readers — Spark's vectorized reader vs. Comet's native scan. So the
contamination does not cancel. The `relative` column is a ratio of sums,
`(scan_spark + expr_spark + c2r_spark) / (scan_comet + expr_comet +
c2r_comet)`, which for any expression cheaper than the scan converges to the
*scan* ratio rather than the expression ratio.
**Result transfer + columnar-to-row.** The `noop` sink writer is a
`DataWriter[InternalRow]`, so a `ColumnarToRowExec` sits on top of every plan.
This is known and tolerated — `CometPlanChecker.findFirstNonCometOperator`
explicitly whitelists `ColumnarToRowExec`, `CometColumnarToRowExec`,
`CometNativeColumnarToRowExec` and `WholeStageCodegenExec`. Since
`spark.comet.exec.columnarToRow.native.enabled` defaults to `false`, the Comet
arm pays Arrow FFI export to JVM `CometVector` to `UnsafeRow`, while Spark
emits rows directly from a single fused codegen stage.
The second cost tracks the *output* type of the expression, so contamination
varies a lot within a single suite:
| Output type | Affected suites | c2r cost |
| --- | --- | --- |
| int / bool / long | comparison, hash, `length`/`ascii`/`instr`,
`like`/`rlike`, `unix_timestamp`, cast-to-numeric | small |
| String | all cast-to-string suites, most of
`CometStringExpressionBenchmark`, `to_json`, `to_csv`, `get_json_object`,
string-result `CASE WHEN` | large |
| struct / array | `from_json`, `sort_array`, `make_interval` | dominant |
So cross-expression rankings within a single results file are distorted too,
not only the absolute numbers.
Worst case is `CometArrayExpressionBenchmark`: 4M rows of 16- and 32-element
int arrays through `.noop()`, where building `UnsafeArrayData` per row very
likely costs more than `sort_array` itself. The suite's own
`element_at(sort_array(x), 1)` variant collapses the output to a single int,
and the gap between those two rows is essentially the c2r cost.
`CometAggregateExpressionBenchmark` is the one suite that is clean on this
axis, because aggregate queries emit one row per group and nothing crosses the
boundary.
## Problem 2: fairness bugs
**ConstantFolding is excluded only for the Comet arm.** `cometExecConfigs`
sets `spark.sql.optimizer.excludedRules` to `ConstantFolding`; the Spark case
runs with folding enabled. Concrete casualty: `select space(2) from
parquetV1Table` in `CometStringExpressionBenchmark` — Spark folds it to a
literal and does no per-row work while Comet evaluates it per row, so that row
reports a fabricated Comet regression. The conf is also *set* rather than
appended, clobbering any pre-existing exclusions.
**`CometPredicateExpressionBenchmark` measures Parquet filter pushdown, not
`In`.** The query is `select * from parquetV1Table where c1 in
('positive','zero')` over a column with three distinct values. Both engines
push that into the reader, so the `In` expression may never be evaluated as an
expression at all. `CometComparisonExpressionBenchmark` gets this right by
putting predicates in the SELECT list instead.
**Two suites use a cardinality too low to measure anything.**
`CometStringExpressionBenchmark` and `CometRegExpBenchmark` both use 1024 rows
— a single batch. Spark's `Benchmark` runs 2s of warmup and then at least 2s of
timed iterations, each iteration being a full Spark job, so these are measuring
job submission and Comet's per-query native plan construction rather than
expression throughput. The `relative` column will read approximately 1.0x
regardless of the expression. That is 31 string expressions and 5 regex
patterns currently producing noise.
## Problem 3: uncontrolled inputs
**Parquet dictionary encoding is left at default in every expression
suite**, while `CometColumnarToRowBenchmark`, `CometExecBenchmark` and
`CometShuffleBenchmark` all set `parquet.enable.dictionary=false` for exactly
this reason. The expression suites diverged from an existing convention in the
same package. This matters because kernels with a dictionary fast path only
compute over distinct values, and distinct counts happen to range from 3
(`CometPredicateExpressionBenchmark`) to 1M
(`CometConditionalExpressionBenchmark.c4`) across suites — so some expressions
get a large unrepresentative win by accident. `CometArithmeticBenchmark` is the
only suite that varies this, and only indirectly via `useDictionary` reducing
distinct value counts.
**Input data is non-deterministic.** `runBenchmarkWithTable` builds the base
table from an unseeded `Random.nextLong()`, so committed results are not
reproducible. This is not merely cosmetic: `CometStringExpressionBenchmark`
derives `REPEAT(CAST(value AS STRING), 10)` from it, giving 10-200 character
strings that vary run to run, which means `lpad(c1, 150, 'x')` and `rpad` are
sometimes padding and sometimes truncating.
**Fallback warnings never reach the results file.** `runExpressionBenchmark`
reports a non-native plan with `println`, which goes to the console, whereas
`Benchmark(output = output)` writes to the `.txt`. Anyone reading a committed
results file cannot tell that a row labelled "Comet" actually ran on Spark.
## Proposed work
Rebuild the harness so that scan and result-transfer costs are either
removed from the timed region or explicitly reported, then migrate the suites
onto it.
- [ ] Add a no-expression baseline case to the shared harness. Take a
`baselineQuery` (default: project the raw input columns) and emit `Spark
(baseline)` / `Comet (baseline)` rows alongside every measurement, so the
scan-plus-c2r floor is visible in every committed results file. This is
mechanical and immediately stops people misreading a 1.1x as an expression win.
- [ ] Add an aggregate-sink measurement mode for expressions whose output is
a string or a complex type, e.g. `SELECT sum(xxhash64(<expr>)) FROM t`, so the
result never crosses the JNI or row boundary. Comet aggregates natively, output
is one row, c2r goes to zero. The added hash per row is small relative to
`UnsafeRow`/`UnsafeArrayData` construction and is identical in both arms.
- [ ] Apply `excludedRules` to both arms rather than only the Comet arm, and
append to the existing value instead of overwriting it.
- [ ] Assert that the benchmarked expression is actually present in the
Spark baseline plan. This catches the `space(2)` class of bug automatically.
- [ ] Raise `CometStringExpressionBenchmark` and `CometRegExpBenchmark` to
at least 1M rows.
- [ ] Move the `In` predicate in `CometPredicateExpressionBenchmark` into
the SELECT list, or disable `spark.sql.parquet.filterPushdown` for that suite.
- [ ] Set `parquet.enable.dictionary` explicitly in
`CometBenchmarkBase.getSparkSession` to match the three other suites in the
package, and make dictionary-encoded vs. plain a deliberate variant where it is
interesting rather than an accident of the generated data.
- [ ] Seed the `Random` in `runBenchmarkWithTable` so results are
reproducible.
- [ ] Route the not-fully-native warning to `output` so it lands in the
results file, and consider failing the run outright.
- [ ] Migrate the existing suites onto the new harness and regenerate the
committed results.
`CometRegExpBenchmark` is worth calling out as already having close to the
right shape: its `Comet (Scan)` case isolates the exec contribution far better
than a plain Spark-vs-Comet pair does. It just runs at 1024 rows.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]