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:
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:
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.
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:
- 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.
- 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:
- 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.
- Is there a seam I have missed that would let Comet intercept either path without a Spark change?
- If we go upstream, narrow patches at the two call sites, or a general
executeArrow on SparkPlan?
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?
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:Two conversions where the ideal number is zero.
The cause is that the Arrow collect paths take the row RDD unconditionally, with no
supportsColumnarcheck. InDataset.toArrowBatchRdd:(
sql/core/.../classic/Dataset.scala:2343on 4.1.3, same shape atDataset.scala:4256on 3.4.3.)collectAsArrowToPythonandtoPythonIteratorboth go through it, so this coverstoPandas(),toArrow(), andcollect()whenspark.sql.execution.arrow.pyspark.enabledis set.Spark Connect has the same problem via a different code path — it does not use
toArrowBatchRddat all:(
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, orcollectAsArrowToPythonanywhere. 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 —
EliminateRedundantTransitionsstrips exactly this round-trip for the UDF path, rewritingMapInArrow/MapInPandasover a Comet C2R intoCometMapInBatchExec(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.
This needs no Comet-specific code in Spark. Any plugin whose
executeColumnaroutput matchesArrowUtils.toArrowSchemacan 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.
toArrowBatchRddisprivate[sql]onDataset, the Connect handler is internal to the Connect module, andSparkSessionExtensionsreaches planner and columnar rules but not either of these.Since Connect and
toArrowBatchRddare separate paths, there is a design question about whether to patch both call sites or introduce something more general, e.g. an optionalexecuteArrowonSparkPlanthat 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:
CometDataset.toArrowBatches(df)or similar) that users call explicitly. Cheap, but only helps people who rewrite their code, and does nothing for Connect.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'sArrowUtils.toArrowTypefor every type it supports, with one exception: Comet has nolargeVarTypesvariant and always emitsUtf8/Binarywith 4-byte offsets. Underspark.sql.execution.arrow.useLargeVarTypes=truea client would getstring/binarywhere it asked forlarge_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 expectArrowBatchStreamWriterframing: schema once, then N uncompressed record batches. This needs a separate serializer rather than a reuse ofserializeBatches.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 byspark.comet.batchSize(8192 by default) whilespark.sql.execution.arrow.maxRecordsPerBatchdefaults 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'stoArrowFieldhas no equivalent.Open questions:
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.executeArrowonSparkPlan?spark.sql.execution.arrow.pyspark.enableddefaults 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?