apache / apache/datafusion-comet

[EPIC] Planner, serde, and optimizer rule performance audit

Open
#5,199 1 comment 0 reactions 1 assignee Claimed by @andygrove View on GitHub
enhancement EPIC performance
Dominant language
Scala
Stars
1.3k
Forks
373
Avg merge
2d 6h
Merged PRs (30d)
190

Description

## What / Why

We spend a lot of effort profiling and optimizing native *execution*, but comparatively little on the JVM-side **planner**: `CometScanRule` / `CometExecRule`, the `QueryPlanSerde` expression and operator serde, and the post-columnar rules. That cost is paid on the driver for every query, and again for **every query stage** under AQE (see the rule injection points documented in `CometSparkSessionExtensions`), so it shows up as latency on short queries and on workloads with many small stages.

This EPIC collects findings from a first pass over that code. Measurements below are from a local probe on a modest 2-table join + aggregate query (778-char SQL, 90 `Expr` nodes in the serialized plan), Spark 4.1 / JDK 17, debug native build:

```
executedPlan w/ Comet = 23.53 ms
executedPlan w/o Comet = 19.05 ms
CometScanRule = 0.23 ms
CometExecRule = 1.13 ms
total serialized plan bytes = 58445
query_context serialized bytes = 56227 (96.2% of plan)
```

So the rules themselves are a modest ~4.5 ms of overhead here, but the **serialized plan is 32x larger than it needs to be**, which is paid per task on every executor.

## Findings

### 1. Serialized plan duplicates the full SQL text on every expression

Tracked separately in https://github.com/apache/datafusion-comet/issues/5200 — 96% of plan bytes, 32x plan-size inflation, 2.4x per-task protobuf round-trip cost. This is the largest item by a wide margin.

### 2. Uncached `Class.getMethod` throughout the Iceberg reflection paths

Likely the largest planner-time cost on Iceberg workloads. `IcebergReflection.extractFileLocation` resolves `location()` on *every* file and detects older Iceberg by **catching** `NoSuchMethodException`, so on `path()`-only versions every call constructs an exception with a stack trace. It is called once per file scan task (`CometIcebergNativeScan.scala`, in the `tasksCollection.asScala.foreach { task => ... }` loop) and once per delete file. In the same loop, `extractDeleteFilesList` re-resolves `keyMetadata()` per task, and `serializePartitionInfo` resolves `spec`, `partitionType`, `fields`, `fieldId` and `get` per task *and* per partition field. `CometScanRule.validateIcebergFileScanTasks` does `getMethod("transform")` + `setAccessible` inside its per-task loop. `IcebergReflection.getFieldInfo` resolves three methods per column.

All of these resolve to a fixed handful of `Method` objects. `Class.getMethod` does a linear search and allocates a fresh `Method` copy on each call, so a scan with tens of thousands of file tasks does O(100k) reflective lookups per planning pass — repeated per AQE stage. There is already a `// TODO` acknowledging this in `IcebergReflection.getFileFormat`. Hoisting the lookups into lazy vals / passing already-resolved `Method`s down is mechanical.

### 3. `tagUnsafePartialAggregates` serializes aggregates only to throw the result away

`CometExecRule.canAggregateBeConverted` calls `QueryPlanSerde.exprToProto` / `aggExprToProto` on grouping, aggregate and result expressions purely to get a `Boolean`. The protos are discarded and then rebuilt during the real `transform` pass, and the probe also burns `exprIdCounter` entries and tags fallback reasons speculatively.

The guards short-circuit, so this only fires for `Final` aggregates whose functions are not mixed-execution-safe, and for multi-stage `collect_list` / `collect_set` — it is not on the path of every query. But when it does fire it is a full double serialization, once per AQE stage. The method's own docstring already notes that a shared predicate helper would be preferable.

### 4. `DecimalPrecision.promote` is a second full traversal of every expression tree

`QueryPlanSerde.exprToProto` runs `expr.transformUp` over the whole tree before the serde walks it again. A cheap guard (does the tree contain a decimal `BinaryArithmetic` at all?) or folding the promotion into the serde walk would remove one of the two traversals. Small per call, but it is called once per projection / filter / aggregate expression.

### 5. `EliminateRedundantTransitions.hasCometNativeChild` is quadratic

It does a full `op.exists(...)` subtree scan at each `ColumnarToRowExec` encountered during a `transformUp`, so cost is quadratic in the number of C2R nodes. Bottom-up memoization is straightforward.

### 6. `PlanDataInjector.injectPlanData` rebuilds the whole operator tree per task

It recurses into every subtree and rebuilds every operator via builders even when only one leaf scan is actually injected. Returning the original `op` unchanged when no descendant needs injection would leave only the path to the scan to rebuild. This runs per task on the executor.

### 7. Micro-optimizations

- `CometConf.getBooleanConf` builds a config-key string, does a map lookup and allocates a `toLowerCase` copy — twice per expression (`isExprEnabled` + `isExprAllowIncompat`).
- `ExpressionRegistry::can_handle` followed by `create_expr` computes `get_expression_type` and hashes twice where a single `get` would do (`native/core/src/execution/planner/expression_registry.rs`).
- `CometScanRule` allocates the Iceberg metadata-table name `Set` on every rule invocation.

## Not investigated yet

- Whether the Comet rules can cheaply detect "already converted" on AQE stage re-entry and skip more work.
- Cost of `CometScanRule` on wide schemas (`isSchemaSupported` walks the full schema per scan, per rule invocation).
- Whether a committed planner benchmark belongs in the repo so these are tracked over time. The probe used for the numbers above was throwaway; turning it into a proper benchmark is worth doing.

## Reproducing

The numbers above came from a temporary ScalaTest suite that plans a query repeatedly with the Comet configs toggled (interleaving the configs across rounds so JIT warm-up does not bias whichever config is measured first — this matters a lot here), reads `CometNativeExec.serializedPlanOpt`, and walks the parsed `Operator` proto counting `Expr` nodes and `query_context` bytes.

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.