Skip to content
Merged
1 change: 1 addition & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -523,6 +523,7 @@ jobs:
org.apache.comet.exec.CometGenerateExecSuite
org.apache.comet.exec.CometWindowExecSuite
org.apache.comet.exec.CometJoinSuite
org.apache.comet.exec.CometTypedDatasetSuite
org.apache.spark.sql.comet.CometMapInBatchSuite
org.apache.spark.sql.execution.python.CometArrowPythonRunnerSuite
org.apache.comet.CometNativeSuite
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -227,6 +227,7 @@ jobs:
org.apache.comet.exec.CometGenerateExecSuite
org.apache.comet.exec.CometWindowExecSuite
org.apache.comet.exec.CometJoinSuite
org.apache.comet.exec.CometTypedDatasetSuite
org.apache.spark.sql.comet.CometMapInBatchSuite
org.apache.spark.sql.execution.python.CometArrowPythonRunnerSuite
org.apache.comet.CometNativeSuite
Expand Down
21 changes: 15 additions & 6 deletions docs/source/contributor-guide/adding_a_new_operator.md
Original file line number Diff line number Diff line change
Expand Up @@ -139,12 +139,21 @@ and leaves the stage itself in place.
`AppendColumnsExec`, `AppendColumnsWithObjectExec`, `MapGroupsExec`, and `CoGroupExec`. They
convert rows to JVM objects, run an arbitrary user function on those objects, or convert them back,
and most of them pass the objects to the next operator as an `ObjectType` column, which has no
Arrow representation. None of this can run natively. Per-row operators can still stay inside a
Comet plan: `MapElementsExec`, the operator behind `Dataset.map`, generates its call to the user
function as a Catalyst `Invoke` expression, so the deserializer, the call, and the serializer could
run together as one projection in the JVM codegen dispatcher. `mapPartitions`, `mapGroups`, and
`cogroup` pass the user function an iterator or a whole group, so there is no per-row expression to
build. A typed `filter` is planned as an ordinary `FilterExec`, not as one of these operators.
Arrow representation. None of this can run natively, so these operators stay on Spark. Every typed
operation ends in `SerializeFromObjectExec`, though, whose output is ordinary rows. With
`spark.comet.convert.typedDataset.enabled`, `CometExecRule` puts a `CometSparkToColumnarExec` above
it, so the operators above the typed operation can run natively. Spark inserts no columnar
transitions below a `RowToColumnarTransition`, so the rule inserts them for the typed operation's
own operators itself. Spark computes a typed operation's rows one at a time, as they are read, while
the conversion fills a whole Arrow batch first. So the rule leaves the output unconverted where a
limit, a `mapPartitions` function, or code reading `Dataset.rdd` could stop reading it early, unless
an operator that reads all of its input first, such as an exchange, a sort, or a hash aggregate,
sits in between. Fusing the deserializer, the `Invoke` that calls the user function, and the
serializer of `Dataset.map` into one projection in the JVM codegen dispatcher was tried in
[#5714](https://github.com/apache/datafusion-comet/pull/5714) and dropped. The dispatcher only calls
into Spark's own classes, and the conversion gets nearly the same speedup for `map` while also
covering the operations that pass the user function an iterator or a whole group. A typed `filter`
is planned as an ordinary `FilterExec`, not as one of these operators.

**Driver-side commands.** `ExecutedCommandExec` runs a `RunnableCommand`, such as DDL or `SET`, on
the driver, so there is no data path for Comet to accelerate.
Expand Down
6 changes: 6 additions & 0 deletions docs/source/contributor-guide/native_shuffle.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,12 @@ Native shuffle (`CometExchange`) is selected when all of the following condition
Spark's `mapsort` normalization makes physical entry order irrelevant. A collated string at any
depth still disqualifies the key. The config defaults to `false` pending measurement of the
nested hashing paths, so by default a complex hash key falls back to JVM shuffle.
- A hash key that is or contains a decimal wider than 18 digits stays on JVM shuffle when the
shuffle's stage starts at a typed `Dataset` conversion
(`spark.comet.convert.typedDataset.enabled`) and the shuffle has more than one partition.
Native shuffle hashes such decimals differently from Spark
([#5994](https://github.com/apache/datafusion-comet/issues/5994)). Without the conversion
this shuffle would have used JVM shuffle, and a join partner may still use it.

## Architecture

Expand Down
4 changes: 4 additions & 0 deletions docs/source/user-guide/latest/datasources.md
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,10 @@ string collations remain unsupported at this conversion boundary. Source default
This includes row-backed `ExistingRDD` inputs when `spark.comet.convert.rdd.enabled=true`. Spark
still produces the RDD rows; conversion lets eligible downstream operators execute in Comet.

The same types apply to the output of typed `Dataset` operations, such as `map`, which Comet
converts when `spark.comet.convert.typedDataset.enabled=true`. A column of any other type keeps
the operators above the typed operation on Spark.
Comment on lines +83 to +85

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The "Other Spark inputs" list above names each spark.comet.convert.* conversion of a Spark operator's output, and this one is missing from it. Could it get a bullet there, for example spark.comet.convert.typedDataset.enabled: the output of typed Dataset operations such as map and mapPartitions? The sentence about column types can stay here. Someone looking for what Comet can convert will not find the new config otherwise.


## Data Catalogs

### Apache Iceberg
Expand Down
2 changes: 1 addition & 1 deletion docs/source/user-guide/latest/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ omitted from the tables below and may be reconsidered based on demand:
- **Structured Streaming operators** (`StateStoreSaveExec`, `StateStoreRestoreExec`, `StreamingSymmetricHashJoinExec`, and similar): Comet targets batch execution.
- **Cartesian / cross joins** (`CartesianProductExec`): rare and expensive, with little acceleration benefit.
- **Pickled (non-Arrow) Python UDFs** (`BatchEvalPythonExec`): Comet accelerates Arrow-based Python UDFs only ([#4234](https://github.com/apache/datafusion-comet/pull/4234)).
- **Typed Dataset operators** (`DeserializeToObjectExec`, `SerializeFromObjectExec`, `MapElementsExec`, `MapPartitionsExec`, `MapGroupsExec`, `CoGroupExec`, `AppendColumnsExec`, and similar): produced by `map`, `mapPartitions`, `groupByKey`, `cogroup`, and other typed `Dataset` transformations. They exist to run user JVM functions on JVM objects, which Comet cannot do natively.
- **Typed Dataset operators** (`DeserializeToObjectExec`, `SerializeFromObjectExec`, `MapElementsExec`, `MapPartitionsExec`, `MapGroupsExec`, `CoGroupExec`, `AppendColumnsExec`, and similar): produced by `map`, `mapPartitions`, `groupByKey`, `cogroup`, and other typed `Dataset` transformations. They exist to run user JVM functions on JVM objects, which Comet cannot do natively. The operators above them can still run natively: set `spark.comet.convert.typedDataset.enabled=true` and Comet converts the output of a typed operation to Arrow. This is disabled by default because it can be slower than Spark when the operators above do little work, such as an aggregate over a few groups, which Spark compiles together with the typed operation into one loop.

## Scans

Expand Down
12 changes: 12 additions & 0 deletions spark/src/main/scala/org/apache/comet/CometConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -248,6 +248,18 @@ object CometConf extends ShimCometConf {
.booleanConf
.createWithDefault(false)

val COMET_CONVERT_FROM_TYPED_DATASET_ENABLED: ConfigEntry[Boolean] =
conf("spark.comet.convert.typedDataset.enabled")
.category(CATEGORY_EXEC)
.doc("When enabled, the output of typed Dataset operations, such as `map`, `flatMap`, " +
"`mapPartitions` and `groupByKey(...).mapGroups`, will be converted to Arrow format so " +
"that the operators above them can run natively. The user function still runs in " +
"Spark. This pays off when the operators above do enough work, such as an " +
"aggregation over many groups, and can be slower when they are cheap, such as an " +
"aggregation over a few groups after a selective filter.")
.booleanConf
.createWithDefault(false)

val COMET_EXEC_ENABLED: ConfigEntry[Boolean] = conf(s"$COMET_EXEC_CONFIG_PREFIX.enabled")
.category(CATEGORY_EXEC)
.doc(
Expand Down
113 changes: 111 additions & 2 deletions spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ import scala.collection.mutable.ListBuffer
import scala.jdk.CollectionConverters._

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.catalyst.expressions.{Divide, DoubleLiteral, EqualNullSafe, EqualTo, Expression, FloatLiteral, GreaterThan, GreaterThanOrEqual, KnownFloatingPointNormalized, LessThan, LessThanOrEqual, NamedExpression, Remainder}
import org.apache.spark.sql.catalyst.expressions.{Divide, DoubleLiteral, EqualNullSafe, EqualTo, Expression, FloatLiteral, GreaterThan, GreaterThanOrEqual, KnownFloatingPointNormalized, LessThan, LessThanOrEqual, NamedExpression, Remainder, SortOrder}
import org.apache.spark.sql.catalyst.expressions.aggregate.{AggregateMode, Final, Partial, PartialMerge}
import org.apache.spark.sql.catalyst.optimizer.NormalizeNaNAndZero
import org.apache.spark.sql.catalyst.rules.Rule
Expand All @@ -49,7 +49,7 @@ import org.apache.spark.sql.execution.datasources.v2.{BatchScanExec, V2CommandEx
import org.apache.spark.sql.execution.datasources.v2.csv.CSVScan
import org.apache.spark.sql.execution.datasources.v2.json.JsonScan
import org.apache.spark.sql.execution.datasources.v2.parquet.ParquetScan
import org.apache.spark.sql.execution.exchange.{BroadcastExchangeExec, BroadcastExchangeLike, ReusedExchangeExec, ShuffleExchangeExec, ShuffleExchangeLike}
import org.apache.spark.sql.execution.exchange.{BroadcastExchangeExec, BroadcastExchangeLike, Exchange, ReusedExchangeExec, ShuffleExchangeExec, ShuffleExchangeLike}
import org.apache.spark.sql.execution.joins.{BroadcastHashJoinExec, BroadcastNestedLoopJoinExec, ShuffledHashJoinExec, SortMergeJoinExec}
import org.apache.spark.sql.execution.window.WindowExec
import org.apache.spark.sql.internal.SQLConf
Expand Down Expand Up @@ -180,6 +180,13 @@ object CometExecRule {

/** Keys of the `spark.comet.convert` configs whose deprecated alternative has been warned. */
private[rules] val warnedDeprecatedConversions = ConcurrentHashMap.newKeySet[String]()

/**
* Tag set on a `SerializeFromObjectExec` whose output an operator above it can stop reading
* early, naming that operator. See `tagPartiallyReadTypedDatasetOutputs`.
*/
private val TYPED_DATASET_PARTIAL_READER: TreeNodeTag[String] =
TreeNodeTag[String]("comet.typedDatasetPartialReader")
}

/**
Expand Down Expand Up @@ -387,6 +394,10 @@ case class CometExecRule(session: SparkSession)
*/
// spotless:on
private def transform(plan: SparkPlan): SparkPlan = {
if (CometConf.COMET_CONVERT_FROM_TYPED_DATASET_ENABLED.get(conf)) {
tagPartiallyReadTypedDatasetOutputs(plan)
}

def convertNode(op: SparkPlan): SparkPlan = op match {
// Scan marker produced by an optional, out-of-tree scan contrib (e.g. contrib/delta).
// Matched by trait (no compile-time dependency on the contrib) and present only when that
Expand Down Expand Up @@ -491,6 +502,23 @@ case class CometExecRule(session: SparkSession)
case op if shouldApplySparkToColumnar(conf, op) =>
convertToComet(op, CometSparkToColumnarExec).getOrElse(op)

// Typed Dataset operations (`map`, `flatMap`, `mapPartitions`, `mapGroups`, ...) pass JVM
// objects between their operators, so those stay on Spark. Each of them ends in
// `SerializeFromObjectExec`, though, whose output is ordinary rows, and converting those to
// Arrow lets the operators above the typed operation run natively.
case op: SerializeFromObjectExec
if CometConf.COMET_CONVERT_FROM_TYPED_DATASET_ENABLED.get(conf) =>
op.getTagValue(CometExecRule.TYPED_DATASET_PARTIAL_READER) match {
case Some(reader) =>
withFallbackReason(
op,
"Comet does not convert the output of a typed Dataset operation when " +
s"$reader can stop reading it early, because filling an Arrow batch would " +
"run the user function on rows that Spark never reaches")
case None =>
convertTypedDatasetOutput(op)
}

// Spark 4.0+: replace only the per-task write, leaving DataWritingCommandExec - and
// therefore Spark's commit protocol, stats trackers and SaveMode handling - in place.
// `V1WritesUtils.getWriteFilesOpt` matches the `WriteFilesExecBase` trait there, which is
Expand Down Expand Up @@ -610,6 +638,10 @@ case class CometExecRule(session: SparkSession)
// Some execs should never be replaced. We include
// these cases specially here so we do not add a misleading 'info' message.
op
case _: ColumnarToRowTransition =>
// A transition does no work of its own. This rule only meets one that
// `convertTypedDatasetOutput` inserted on an earlier pass over the same plan.
op
case _: WriteFilesExec =>
// The write is converted at the enclosing DataWritingCommandExec above: on Spark 3.x
// by replacing the whole command, on 4.0+ by converting this child from there.
Expand Down Expand Up @@ -1274,6 +1306,83 @@ case class CometExecRule(session: SparkSession)
private def hasEnabledHandler(op: SparkPlan): Boolean =
allExecs.get(op.getClass).exists(_.enabledConfig.forall(_.get(op.conf)))

/**
* Tags each `SerializeFromObjectExec` whose output an operator above it can stop reading early
* with that operator, so `transform` does not convert it. Spark computes the rows of a typed
* Dataset operation one at a time, as the operator above reads them, while the conversion fills
* a whole Arrow batch first. Below a limit, a `mapPartitions` function such as `_.take(1)`, or
* code reading `Dataset.rdd`, the conversion would run the user function on rows Spark never
* reaches, and a function that throws on one of them would fail a query that succeeds in Spark.
* An operator that reads all of its input before it returns a row ends the search, since Spark
* computes every row below it anyway: an exchange, a sort, a hash aggregate, or a top-k over
* input that is not already sorted.
*
* Conversion is bottom-up, so this runs first. TreeNode tags survive the child copies made
* during transformUp, while an identity set would not.
*/
private def tagPartiallyReadTypedDatasetOutputs(plan: SparkPlan): Unit = {
def visit(op: SparkPlan, partialReader: Option[String]): Unit = {
val childReader = op match {
case serialize: SerializeFromObjectExec =>
partialReader.foreach(
serialize.setTagValue(CometExecRule.TYPED_DATASET_PARTIAL_READER, _))
partialReader
case _: CollectLimitExec | _: LocalLimitExec | _: GlobalLimitExec => Some("a limit")
// A top-k reads only its first rows when its input is already sorted.
case topK: TakeOrderedAndProjectExec
if SortOrder.orderingSatisfies(topK.child.outputOrdering, topK.sortOrder) =>
Some("a limit")
Comment on lines +1331 to +1334

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can this arm ever tag a SerializeFromObjectExec? Typed operators report no outputOrdering (SerializeFromObjectExec does not override it), so an ordering that satisfies the top-k comes from a SortExec, and the SortExec arm below resets the search before it reaches the serializer. I could not build a plan where this arm changes a tag, and no test reaches it. If there is none, dropping the arm, its comment and the SortOrder import leaves TakeOrderedAndProjectExec as a plain barrier. If there is one, a test would pin it.

case _: MapPartitionsExec => Some("a mapPartitions function")
case _: Exchange | _: SortExec | _: HashAggregateExec | _: ObjectHashAggregateExec |
_: TakeOrderedAndProjectExec =>
None
case _ => partialReader
}
op.children.foreach(visit(_, childReader))
}
// `Dataset.rdd` reads the objects the plan produces, and the RDD's own code decides how many
// of them to read, as `take(1)` does. The plan's root is the `DeserializeToObjectExec` that
// `Dataset.rdd` adds. When the Dataset ends in a typed operation such as `map`, Spark's
// `EliminateSerialization` drops that deserializer together with the operation's serializer,
// so the root is the operation itself, which produces objects too, or a typed filter or a
// project over it. A Dataset's own plan ends in rows, so no other plan has such a root.
def producesObjects(op: SparkPlan): Boolean = op match {
case _: ObjectProducerExec => true
case _: FilterExec | _: ProjectExec => producesObjects(op.children.head)
case _ => false
}
visit(plan, if (producesObjects(plan)) Some("code reading Dataset.rdd") else None)
}

/**
* Converts the rows a typed Dataset operation produces to Arrow, so the operators above it can
* run natively. See [[CometConf.COMET_CONVERT_FROM_TYPED_DATASET_ENABLED]].
*
* Spark inserts the columnar transitions after this rule, but it does not look below a
* `RowToColumnarTransition` such as `CometSparkToColumnarExec`. That is harmless above a leaf.
* Here the typed operation's own operators sit below the conversion, and without a transition
* they would read a Comet child through `CometExec.doExecute`, Spark's interpreted
* columnar-to-row path. So the subtree gets its transitions now, from Spark's own rule, and
* `EliminateRedundantTransitions` later replaces each one over a Comet child with Comet's own.
* Spark's rule leaves existing transitions alone, which matters because this rule runs over the
* same plan twice under AQE.
*/
private def convertTypedDatasetOutput(op: SerializeFromObjectExec): SparkPlan = {
val unsupported = op.output.filterNot(a =>
CometSparkToColumnarExec.isTypeSupported(a.dataType, a.name, ListBuffer.empty))
if (unsupported.nonEmpty) {
withFallbackReason(
op,
"Comet cannot convert the output of a typed Dataset operation to Arrow because it does " +
"not support the type of these columns: " +
unsupported.map(a => s"${a.name}: ${a.dataType.simpleString}").mkString(", "))
} else {
val withTransitions =
ApplyColumnarRulesAndInsertTransitions(Seq.empty, outputsColumnar = false).apply(op)
convertToComet(withTransitions, CometSparkToColumnarExec).getOrElse(withTransitions)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Prevent conversion from creating incompatible wide-decimal shuffles. With this feature enabled, a supported typed input switches to native shuffle while another typed input containing a retained array<int> column stays on JVM shuffle. Joining their decimal(38,18) keys with AQE disabled then silently loses matching rows: my reproduction returns 9 instead of 100. Conversion-disabled Comet returns all 100. This newly exposes the existing decimal hash difference as incorrect query results. Keep affected exchanges on JVM shuffle until native wide-decimal hashing matches Spark, and cover this mixed-path join.

Evidence: Reproduced on Spark 4.1.3/JDK 17 in CometTestBase with spark.sql.adaptive.enabled=false, spark.sql.autoBroadcastJoinThreshold=-1, spark.sql.shuffle.partitions=10, and spark.comet.shuffle.mode=auto. Define case class L(k: java.math.BigDecimal, v: Long) and case class R(k: java.math.BigDecimal, xs: Seq[Int]). Build l = spark.range(0,100,1,2).map(i => L(new java.math.BigDecimal(i), i)).alias("l") and r = spark.range(0,100,1,2).map(i => R(new java.math.BigDecimal(i), Seq(i.toInt))).alias("r"). Collect l.join(r, col("l.k") === col("r.k")).select(col("l.v"), col("r.xs")). Spark and conversion-disabled Comet return 100 rows. Conversion-enabled Comet returns 9. The executed plan contains left CometNativeShuffle and right CometColumnarShuffle. Setting shuffle mode to jvm restores 100 rows. Spark hashes wide decimals using unscaledValue().toByteArray; native hash_array_decimal! uses fixed-width to_le_bytes().

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the repro. It reproduced exactly here: 9 rows instead of 100, with CometNativeShuffle on the left and CometColumnarShuffle on the right.

Fixed in 1ec0e66. A hash shuffle whose stage starts at a typed Dataset conversion now stays on Comet's columnar shuffle when a key is or contains a decimal wider than 18 digits and the shuffle has more than one partition, so the conversion no longer moves it off the path the other side of the join uses. That's the rule #6005 applies to every native shuffle, limited to the shuffles this feature moves, so it becomes redundant once #6005 lands. Your join is now a test (a join on wide decimal keys with an input that is not converted), which fails without the guard. The rule is also listed in native_shuffle.md.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Preserve short-circuiting when a limit consumes typed output. With spark.comet.convert.typedDataset.enabled=true, map(...).limit(1) fills an Arrow batch before returning its first row. A user function that throws on row 30 therefore fails the query, although Spark and conversion-disabled Comet return the first row successfully. The resulting plan is CometCollectLimit -> CometSparkRowToColumnar -> SerializeFromObject. Please preserve row-level limiting before batching where valid, or decline conversion for this pipeline, and add a regression test.

Evidence: Reproduced at the reviewed head in CometTestBase on Spark 4.1.3/JDK 17, with AQE both false and true: spark.range(0, 100, 1, 1).map { i => if (i == 30L) throw new IllegalArgumentException("unexpected evaluation of row 30"); i + 1L }.toDF().limit(1).collect(). Spark and Comet with typed conversion disabled return [Row(1)]. Enabling conversion throws SparkException caused by that IllegalArgumentException, with RowArrowReader.loadNextBatch in the stack. Spark’s CollectLimitExec.executeCollect calls child.executeTake(limit), whereas the inserted Arrow reader consumes a batch before the limit can stop it. The six-configuration probe failed only in the two conversion-enabled cases.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 2737816, and generalized in c202f51. The conversion leaves a typed operation's output unconverted when a limit above it can stop reading early, unless an operator that reads all of its input, such as an exchange, a sort or a hash aggregate, sits in between. Your query is the test a limit does not evaluate typed Dataset rows beyond the result, with AQE on and off.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A native shuffle placed directly over this conversion reads it through executeColumnar() and ColumnarBatchArrowReader, which closes each batch. RowArrowReader reuses its vectors, and closing a struct vector drops its children, so a struct column fails from the second batch. #6607 hits the same thing and fixes it by reading a CometNativeArrowSource child as an Arrow stream.

I replayed this with the CometExecRule, CometShuffleExchangeExec and CometConf from this head over a local build of main (Spark 4.1.3, spark.comet.convert.typedDataset.enabled=true, spark.comet.batchSize=16):

case class Rec(a: Int, b: String)
case class Out(id: Int, inner: Rec)
spark.range(0, 200, 1, 2).map(i => Out(i.toInt, Rec(i.toInt, s"s$i"))).repartition(col("id")).collect()

This fails with CometNativeException: ... no more field nodes for field a, with AQE on and off. orderBy("id"), the left side of a shuffled join, and a union that feeds a shuffle fail the same way. The plan is CometExchange ... CometNativeShuffle over CometSparkRowToColumnar over SerializeFromObject. The same repartition returns all 200 rows with one batch, with flat columns, with an array<string> column, or with the conversion off.

No test here puts a native shuffle directly over the conversion. The struct test reads only id downstream, so Spark's ObjectSerializerPruning removes inner and tags from the serializer before the conversion sees them (a plain explain of that query shows SerializeFromObject [... AS id]). Could you add a struct case with a shuffle directly above the conversion and a small spark.comet.batchSize? It will fail until #6607 lands, so please either land that first or decline struct columns in convertTypedDatasetOutput until then.

}
}

private def shouldApplySparkToColumnar(conf: SQLConf, op: SparkPlan): Boolean = {
// Only consider converting leaf nodes to columnar currently, so that all the following
// operators can have a chance to be converted to columnar. Leaf operators that output
Expand Down
Loading
Loading