apache / apache/datafusion-comet

Native columnar-to-row conversion is much slower than the JVM implementation for small batches

Open
#5,112 2 comments 0 reactions 0 assignees View on GitHub
area:ffi bug performance priority:medium
Dominant language
Scala
Stars
1.3k
Forks
373
Avg merge
2d 6h
Merged PRs (30d)
190

Description

### Describe the issue

`CometNativeColumnarToRowExec` (enabled by default) can be much slower than the JVM `CometColumnarToRowExec`, but none of our existing benchmarks can see it. Measuring the conversion in isolation (no parquet scan) shows the native converter is 3.7x slower per row at the default 8192-row batch size, and the gap grows quickly as batches get smaller because the native path carries a large fixed cost per batch.

Isolated microbenchmark (in-memory Arrow batches, 4 columns: long, int, double, string; conversion plus value consumption only; Apple M3 Ultra, Spark 4.1):

| batch size | JVM rowIterator + UnsafeProjection | native converter | native/JVM |
|---|---|---|---|
| 8192 | 8.8 ns/row | 32.6 ns/row | 3.7x slower |
| 512 | 9.1 ns/row | 53.7 ns/row | 5.9x slower |
| 32 | 19.5 ns/row | 306.6 ns/row | 15.7x slower |

Solving the three data points for a fixed-plus-linear cost model gives roughly 10us of fixed overhead per batch for the native path versus roughly 0.2us for the JVM path. Note that the JVM baseline here is the interpreted iterator path; the production JVM operator is `CodegenSupport` and fuses into whole-stage codegen, so it is cheaper still.

### Where the native cost goes

1. Per-batch fixed overhead (~10us, the dominant real-world problem): every batch pays Arrow FFI export on the JVM side (`allocateArrowStructs` plus `Data.exportVector` per column, re-exporting the constant schema every time), the JNI crossing, `from_ffi` import per column in Rust, and on the return path `FindClass` for `NativeColumnarToRowInfo`, object construction, and two freshly allocated `int[]` arrays. `from_ffi` uses `ArrayData::new_unchecked`, so validation is not a factor; it is the number of per-batch allocations and crossings.
2. Per-row cost (3.7x even at the ideal batch size): for any schema containing a var-length column the Rust side takes the general path (per-row buffer resize and zero, two passes over the columns, enum dispatch per column per row), and every row then crosses back through `NativeRowIterator.next()`.
3. The per-row `unsafeRow.copy()` added in #3367 to fix the #3308 SIGSEGV costs only ~4 ns/row on narrow rows, so it is secondary, but it grows with row width, allocates twice per row, and permanently negates the zero-copy premise described in the operator's docs.

Two aggravating factors:

- The operator has no `CodegenSupport`, so it is a hard whole-stage-codegen boundary. The JVM implementation fuses into the downstream stage and often never materializes an `UnsafeRow` at all.
- The `convertTime` metric only times the JNI call. The per-row iteration and copy cost lands in whatever operator consumes the rows, so in the Spark UI the slowness is misattributed to downstream operators.

### Why existing benchmarks show parity

`CometColumnarToRowBenchmark` measures end-to-end parquet queries with a `noop()` sink at 70-530 ns/row, where the scan swamps the conversion, and the scan always produces full 8192-row batches, which amortizes the per-batch overhead. It shows parity across every group, and stays at parity even with dictionary-encoded input or a Scala UDF consumer. This also matches the TPC-DS SF1000 parity result in #4440.

The workloads that hit the slow path are the ones that reach C2R with small batches: selective native filters, join outputs, small shuffle partitions, broadcast dimension tables, and AQE-coalesced partitions.

### Possible improvements

- Amortize the per-batch cost: cache the jclass/methodID for `NativeColumnarToRowInfo`, reuse the FFI structs, return offsets and lengths through a preallocated direct buffer, and export the schema once per converter instead of once per batch.
- Add a per-batch dynamic fallback: below some row-count threshold, convert with the JVM path. Both implementations are available at the call site.
- Revisit the blanket `copy()`: transferring ownership of the row buffer to the returned iterator (freed on close, double-buffered in Rust) would restore Spark's normal row-reuse contract without reintroducing the #3308 SIGSEGV.
- Extend `CometColumnarToRowBenchmark` with a small-batch scenario, for example feeding C2R through a selective native filter, so this regression class is visible.

Related: #4440 (questions the operator's value given parity results), #3308, #3367.

### Incidental find

`CometColumnarToRowBenchmark` is currently broken on Spark 4.x defaults: `cast(id as byte)` overflows under ANSI mode at id=128. It needs `spark.sql.ansi.enabled=false` in its session.

Contributor guide

Open the contributing guide

Research direction

Start with CometNativeColumnarToRowExec, CometColumnarToRowExec, NativeRowIterator.next(), and CometColumnarToRowBenchmark; reproduce the isolated small-batch comparison and fix the benchmark's Spark 4.x ANSI-mode failure. The issue lists several possible optimizations but does not select one, so done requires agreeing on an approach and demonstrating improved small-batch conversion without regressing existing benchmark coverage.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, rust, scala
Domain
backend, performance, testing
Issue type
Bug
Difficulty
5/5
Estimated time
Over a week
Activity status
Quiet
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.