apache / apache/datafusion-comet
Arrow-based clients pay a redundant Arrow -> UnsafeRow -> Arrow round-trip
- Dominant language
- Scala
- Stars
- 1.3k
- Forks
- 373
- Avg merge
- 2d 4h
- Merged PRs (30d)
- 198
Description
### What is the problem the feature request solves?
Comet produces Arrow. Arrow-based clients consume Arrow. Between the two, Spark inserts a full round-trip through `UnsafeRow`:
```
Comet Arrow batch -> C2R -> UnsafeRow -> ArrowConverters -> Arrow IPC -> client
```
Two conversions where the ideal number is zero.
The cause is that the Arrow collect paths take the row RDD unconditionally, with no `supportsColumnar` check. In `Dataset.toArrowBatchRdd`:
```scala
private[sql] def toArrowBatchRdd(plan: SparkPlan): RDD[Array[Byte]] = {
...
plan.execute().mapPartitionsInternal { iter =>
val context = TaskContext.get()
ArrowConverters.toBatchIterator(iter, schemaCaptured, maxRecordsPerBatch, timeZoneId, ...)
}
}
```
(`sql/core/.../classic/Dataset.scala:2343` on 4.1.3, same shape at `Dataset.scala:4256` on 3.4.3.) `collectAsArrowToPython` and `toPythonIterator` both go through it, so this covers `toPandas()`, `toArrow()`, and `collect()` when `spark.sql.execution.arrow.pyspark.enabled` is set.
Spark Connect has the same problem via a *different* code path — it does not use `toArrowBatchRdd` at all:
```scala
val rows = dataframe.queryExecution.executedPlan.execute()
...
val batches = rows.mapPartitionsInternal(
SparkConnectStreamHandler.rowToArrowConverter(schema, maxRecordsPerBatch, maxBatchSize, timeZoneId))
```
(`SparkConnectStreamHandler.scala:113`, verified on 3.4.3 — the handler was reorganised in 4.x and I have not checked whether the shape survived, which is worth confirming.) Connect always transfers Arrow, so unlike the PySpark case there is no conf to turn this off.
I grepped the Comet tree and we intercept none of this today: no handling of `toArrowBatchRdd`, `ArrowConverters`, or `collectAsArrowToPython` anywhere. So a Comet user on PySpark-with-Arrow or on Connect pays for the C2R and then pays again to rebuild what Comet already had in the right format.
I think this is worth separating from the existing C2R work. #5119, #5112, and #4440 are all about making the conversion itself cheaper. This is a case where the conversion's output is discarded and re-encoded immediately, so the available win is deleting both halves rather than shaving one.
We already know the trick — `EliminateRedundantTransitions` strips exactly this round-trip for the UDF path, rewriting `MapInArrow`/`MapInPandas` over a Comet C2R into `CometMapInBatchExec` (`EliminateRedundantTransitions.scala:124`). The collect path is the same shape, minus a seam to hook.
### Describe the potential solution
The clean fix is upstream: make the Arrow collect paths columnar-aware.
```scala
if (plan.supportsColumnar) {
plan.executeColumnar().mapPartitionsInternal { /* serialize batches straight to Arrow IPC */ }
} else {
plan.execute().mapPartitionsInternal { ArrowConverters.toBatchIterator(...) }
}
```
This needs no Comet-specific code in Spark. Any plugin whose `executeColumnar` output matches `ArrowUtils.toArrowSchema` can take the fast path, so Gluten and the RAPIDS plugin benefit identically — which should help the case upstream.
Comet cannot do this on its own. `toArrowBatchRdd` is `private[sql]` on `Dataset`, the Connect handler is internal to the Connect module, and `SparkSessionExtensions` reaches planner and columnar rules but not either of these.
Since Connect and `toArrowBatchRdd` are separate paths, there is a design question about whether to patch both call sites or introduce something more general, e.g. an optional `executeArrow` on `SparkPlan` that any Arrow consumer can ask for and that falls back to the row path when unimplemented. The second is more work to land but stops this from recurring at the next Arrow consumer.
Interim options, if we want anything before an upstream change:
1. A Comet-side entry point (`CometDataset.toArrowBatches(df)` or similar) that users call explicitly. Cheap, but only helps people who rewrite their code, and does nothing for Connect.
2. Nothing, and treat the upstream change as the deliverable.
I lean towards (2), with the measurement below done first, but I would like other opinions.
### Additional context
Things that need resolving before this is more than an idea.
**Schema equivalence.** Comet's `Utils.toArrowType` (`spark/src/main/scala/org/apache/spark/sql/comet/util/Utils.scala:139`) is field-for-field identical to Spark's `ArrowUtils.toArrowType` for every type it supports, with one exception: Comet has no `largeVarTypes` variant and always emits `Utf8`/`Binary` with 4-byte offsets. Under `spark.sql.execution.arrow.useLargeVarTypes=true` a client would get `string`/`binary` where it asked for `large_string`/`large_binary`. That is the same conflict we already resolve by falling back for mapInArrow (`EliminateRedundantTransitions.scala:113-118`), so the same fallback applies.
**Batch framing and compression.** `Utils.serializeBatches` (`util/Utils.scala:242`) writes one complete Arrow stream per batch, compressed with Spark's codec. That is what the broadcast path wants, but clients expect `ArrowBatchStreamWriter` framing: schema once, then N uncompressed record batches. This needs a separate serializer rather than a reuse of `serializeBatches`.
**Dictionary-encoded vectors.** Comet can emit `CometDictionaryVector`. pyarrow handles dictionary arrays, but they surface in pandas as categoricals, which is a visible behavior change from today. Either decode before serializing, or decide the change is acceptable and document it.
**`maxRecordsPerBatch`.** Comet batches are sized by `spark.comet.batchSize` (8192 by default) while `spark.sql.execution.arrow.maxRecordsPerBatch` defaults to 10000. Do we honor the Arrow conf by re-batching, or document that Comet's batch size wins on this path? Connect additionally caps by estimated byte size (`CONNECT_GRPC_ARROW_MAX_BATCH_SIZE`), which a columnar path would have to respect.
**`errorOnDuplicatedFieldNames`.** Spark threads this into schema construction for the pandas struct handling mode. Comet's `toArrowField` has no equivalent.
Open questions:
1. Does anyone have a measurement of what this actually costs on a realistic `toPandas()` or Connect workload? I have not benchmarked it — the argument so far is structural, from reading the call paths. Worth confirming the win is real, and how large, before proposing anything upstream.
2. Is there a seam I have missed that would let Comet intercept either path without a Spark change?
3. If we go upstream, narrow patches at the two call sites, or a general `executeArrow` on `SparkPlan`?
4. `spark.sql.execution.arrow.pyspark.enabled` defaults to false, so the PySpark half only affects users who opted in. Does anyone have a sense of how common that is among Comet users, relative to Connect?
Contributor guide
Research direction
Start by measuring the redundant conversion on realistic PySpark and Spark Connect workloads. Trace Dataset.toArrowBatchRdd in Dataset.scala and the row conversion in SparkConnectStreamHandler.scala, then compare their schema, framing, batching, compression, and dictionary behavior with Utils.toArrowType and Utils.serializeBatches. Done means the measured benefit and the scope—narrow call-site changes or a general seam—are documented, including required fallbacks.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- scala
- Domain
- data-engineering, distributed-systems
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Active
- Clarity
- Needs clarification
- Newbie friendliness
- 35/100