Skip to content

Arrow-based clients pay a redundant Arrow -> UnsafeRow -> Arrow round-trip #5616

Description

@andygrove

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:

  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?

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions