Skip to content

A Spark operator reading a Comet shuffle through AQE's coalesced read gets Spark's ColumnarToRow instead of Comet's #6610

Description

@andygrove

What is the problem the feature request solves?

EliminateRedundantTransitions replaces a Spark ColumnarToRowExec that reads a Comet operator with Comet's own transition: CometColumnarToRowExec, or CometNativeColumnarToRowExec when spark.comet.exec.columnarToRow.native.enabled is on. It misses a Spark operator that reads a Comet shuffle through AQEShuffleReadExec. AQE inserts that node whenever it coalesces partitions, splits skewed ones or reads shuffle output locally, and coalescing is on by default.

hasCometNativeChild unwraps a QueryStageExec only when it is the transition's direct child, and containsCometPlan treats a stage as a leaf. So AQEShuffleRead(ShuffleQueryStage(CometExchange)) looks as if it holds no Comet operator, and Spark's transition stays.

Reproduced on main (ff468b1) with a typed map over a repartitioned Comet scan, with AQE on:

spark.sql("SELECT _1 AS a, _2 AS b FROM tbl")
  .repartition(col("b"))
  .as[Rec]
  .map(r => Rec(r.a + 1, r.b))

With spark.sql.adaptive.coalescePartitions.enabled=true, the default:

+- *(1) DeserializeToObject newInstance(class Rec), obj#15: Rec
   +- *(1) ColumnarToRow
      +- AQEShuffleRead coalesced
         +- ShuffleQueryStage 0
            +- CometExchange hashpartitioning(b#5, 10), REPARTITION_BY_COL, CometNativeShuffle

With it set to false:

+- *(1) DeserializeToObject newInstance(class Rec), obj#44: Rec
   +- *(1) CometColumnarToRow
      +- ShuffleQueryStage 0
         +- CometExchange hashpartitioning(b#34, 10), REPARTITION_BY_COL, CometNativeShuffle

The cost looks small today:

  • spark.comet.exec.columnarToRow.native.enabled never applies to these transitions.
  • CometColumnarToRowExec is Spark's transition plus the fix for SPARK-50235, which closes each batch's vectors once its rows are read. Spark 3.5.8, 4.0.1 and 4.1.3 already do the equivalent with closeIfFreeable(). Spark 3.4.3's transition has no cleanup, but NativeBatchDecoderIterator closes the previous batch as it advances, so I don't expect a memory difference on any version.
  • ExtendedExplainInfo and RevertNativeForTransitionHeavyStages count Spark's transition and Comet's alike, so plan reports and the transition-heavy stage revert don't change.

Describe the potential solution

Unwrap AQEShuffleReadExec in hasCometNativeChild, so that the check reaches the stage below it. AQEShuffleReadExec reports its stage's supportsColumnar and returns the stage's Comet batches, so Comet's transition can sit on it. A test in EliminateRedundantTransitionsSuite with coalescing on would cover it, and a skew-split read takes the same path.

Additional context

Found while working on #6607, whose tests saw Spark's transition above the coalesced read of a native shuffle.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:shuffleShuffle (JVM and native)enhancementNew feature or request

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions