apache / apache/datafusion-comet

Investigate: scan + shuffle-write slower than Spark on deeply nested schemas

Open
#5,355 1 comment 0 reactions 1 assignee Claimed by @andygrove View on GitHub
area:scan area:shuffle enhancement performance
Dominant language
Scala
Stars
1.3k
Forks
373
Avg merge
2d 4h
Merged PRs (30d)
198

Description

## Background

We have reports that the scan + shuffle-write (map-side) stage is slower than Spark
when the schema contains deeply nested types. This issue collects a static-analysis
pass over both halves of that stage so the candidates can be measured and split into
separate fixes.

**Status: unverified.** Everything below is code reading with exact call sites — no
benchmark numbers yet. The point of the issue is to record the candidates and the
reasoning so they can be confirmed or dismissed one at a time. Some of these may
turn out to be noise next to the actual data copy.

A common thread runs through most of them: work whose cost scales with the *size of
the schema tree* is being done per batch instead of once per plan/writer. That is
invisible on flat schemas and grows with nesting depth and width, which matches the
shape of the reports.

## Scan side

### Candidate A: a trivial nested-type difference re-converts the whole column every batch

`native/core/src/parquet/schema_adapter.rs:589` (and the same gate at `:886`) decide
whether to wrap a column based on:

```rust
if logical_field.data_type() != physical_field.data_type() {
```

That is a whole-tree comparison, and arrow's `Field::eq`
(`arrow-schema-58.4.0/src/field.rs:108`) compares `name`, `data_type`, `nullable`
**and `metadata`**. So a single difference anywhere in the tree wraps the *entire*
top-level nested column in `CometCastColumnExpr`, which then reconverts the whole
subtree on every batch.

Known triggers that differ only cosmetically:

- `PARQUET:field_id` metadata. arrow-rs attaches it per field only when the parquet
schema actually carries IDs (`parquet-58.4.0/src/arrow/schema/mod.rs:513`, guarded
by `basic_info.has_id()`). So this applies to Iceberg-written files and to Spark
with `spark.sql.parquet.fieldId.write.enabled`, **not** to vanilla Spark parquet.
- INT96 timestamps, where `coerce_int96_tz` attaches a UTC timezone.
- Inner field nullability differences at any level. arrow-rs derives physical
nullability from the parquet repetition
(`parquet-58.4.0/src/arrow/schema/mod.rs:497`), so a `REQUIRED` nested field reads
back as `nullable: false`. Whether that matches the logical schema depends on where
the required schema came from (catalog metadata often widens nested fields to
nullable), so this needs checking per table rather than assuming.

Ruled out while writing this up: inner **field naming** is not a trigger. arrow-rs
names list children `element` (`parquet-58.4.0/src/arrow/schema/mod.rs:535`) and
Comet's `Utils.toArrowField` does too — the `element` -> `item` rename lives only in the
codegen FFI path (`spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenOutput.scala:45`),
which is not this path.

The deeper and wider the type, the more chances something mismatches, and the more a
wrap costs once it fires. I have **not** confirmed which trigger (if any) fires on the
reported workloads — that is the first thing to check, and it decides whether this
candidate matters at all.

### Candidate B: the per-batch conversion redoes plan-time work

Once Candidate A fires, `parquet_convert_struct_to_struct` rebuilds two `HashMap`s and
allocates a `String` per field name — per struct, per nesting level, per batch:

- `native/core/src/parquet/parquet_support.rs:290` — `from_id_to_index`
- `native/core/src/parquet/parquet_support.rs:302` — `normalize_name` returns `String`
- `native/core/src/parquet/parquet_support.rs:309-311` — `field_name_to_index_map`

None of this depends on the data. The field-matching plan could be computed once and
reused. This is only worth fixing if Candidate A turns out to be unavoidable in some
cases; if we can stop wrapping cosmetic mismatches, this cost disappears with it.

## Shuffle side — JVM / columnar path (`process_sorted_row_partition`)

### Candidate C: `ShuffleBlockWriter::try_new` is inside the per-batch loop

`native/shuffle/src/spark_unsafe/row.rs:1422` constructs a new `ShuffleBlockWriter`
for every batch, inside the `while current_row < row_num` loop.

`try_new` calls `schema.flattened_fields()` (walks the full nested field tree) and
pre-encodes the entire IPC schema flatbuffer. Its own doc comment says the whole point
is to encode the schema *once* per writer rather than per block:

> For the common case of a schema with no dictionary types, the schema flatbuffer is
> encoded once in `Self::try_new` and written verbatim at the start of every block,
> rather than being re-serialized per block as `StreamWriter::try_new` would do.

Calling it per batch defeats exactly that. Every other call site in the tree builds it
once — e.g. the native path at `native/shuffle/src/shuffle_writer.rs:213`. Cost scales
with schema tree size, so nested schemas pay the most. This looks like the most
clear-cut item here and the fix is contained (hoist above the loop).

Adjacent, same function: `make_batch` (`row.rs:1490`) rebuilds the `Schema` per batch,
and the following `RecordBatch::try_new_with_options` then deep-compares every nested
`DataType` against it.

### Candidate D: field-major struct reads re-walk the row buffers per leaf field

`append_struct_fields_field_major` (`native/shuffle/src/spark_unsafe/row.rs:805`) hoists
type dispatch out of the row loop, but re-runs `read_row_at!` + `get_struct` *inside*
each field's row loop. The traversal became one full sweep over all N `UnsafeRow`
buffers **per leaf field** rather than one sweep total — cache-hostile, and it
compounds with depth via `append_nested_struct_fields_field_major`.

It also allocates three `Vec`s of length `num_rows` per nested struct field
(`row.rs:968-970`).

Worth noting this is a deliberate optimization that traded dispatch cost for traversal
cost; on deep/wide schemas that trade may go the wrong way. Needs measuring against the
row-major predecessor rather than assuming.

## Shuffle side — native path

### Candidate E: `Map` columns are repartitioned row by row

`interleave_record_batch` is the repartition primitive
(`native/shuffle/src/partitioners/partitioned_batch_iterator.rs:111`).

arrow 58.4's `interleave` dispatch has dedicated arms for `Struct` and `List`
(`arrow-select-58.4.0/src/interleave.rs:108-110`) but **none for `Map`** or
`FixedSizeList`, so those fall through to `interleave_fallback` at `:111` →
`MutableArrayData::extend` per contiguous run, which for hash-scattered indices is
effectively per row.

So struct- and list-heavy schemas are fine on this path; map-typed columns are not.
Fix likely belongs upstream in arrow-rs.

### Candidate F (minor): per-batch walk of the whole ArrayData tree

`count_new_buffers` (`native/shuffle/src/partitioners/multi_partition.rs:429`) calls
`to_data()` per column and walks the full nested tree with a `HashSet` insert per
buffer, once per batch. Probably small, listed for completeness — and note the existing
doc comment explains why the cheaper alternatives were rejected, so this one should not
be "optimized" without reading that reasoning first.

## Suggested order

1. Confirm whether Candidate A fires on a real affected plan (count `CometCastColumnExpr`
in the explain output). This decides whether the scan half matters at all.
2. Candidate C — clear-cut and self-contained, independent of the rest.
3. Candidate E — real gap, but only for map-typed columns.
4. Candidates B / D — only once there are numbers justifying them.

There is an existing `shuffle_block_schema_encoding` benchmark group in
`native/shuffle/benches/shuffle_writer.rs` covering flat vs. deeply nested schemas that
can be extended to measure Candidate C directly.

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.