Adding a New Operator#
This guide explains how to add support for a new Spark physical operator in Apache DataFusion Comet.
Overview#
CometExecRule is responsible for replacing Spark operators with Comet operators. There are different approaches to
implementing Comet operators depending on where they execute and how they integrate with the native execution engine.
Types of Comet Operators#
CometExecRule maintains two distinct maps of operators:
1. Native Operators (nativeExecs map)#
These operators run entirely in native Rust code and are the primary way to accelerate Spark workloads. Native
operators are registered in the nativeExecs map in CometExecRule.scala.
Key characteristics of native operators:
They are converted to their corresponding native protobuf representation
They execute as DataFusion operators in the native engine
The
CometOperatorSerdeimplementation handles enable/disable checks, support validation, and protobuf serialization
Examples: ProjectExec, FilterExec, SortExec, HashAggregateExec, SortMergeJoinExec, ExpandExec, WindowExec
2. Sink Operators (sinks map)#
Sink operators serve as entry points (data sources) for native execution blocks. They are registered in the sinks
map in CometExecRule.scala.
Key characteristics of sinks:
They become
ScanExecoperators in the native plan (seeoperator2ProtoinCometExecRule.scala)They can be leaf nodes that feed data into native execution blocks
They are wrapped with
CometScanWrapperorCometSinkPlaceHolderduring plan transformationExamples include operators that bring data from various sources into native execution
Examples: UnionExec, CoalesceExec, CollectLimitExec, TakeOrderedAndProjectExec
Special sinks (not in the sinks map but also treated as sinks):
CometScanExec- File scansCometSparkToColumnarExec- Conversion from Spark row formatShuffleExchangeExec/BroadcastExchangeExec- Exchange operators
3. Comet JVM Operators#
These operators run in the JVM but are part of the Comet execution path. For JVM operators, all checks happen
in CometExecRule rather than using CometOperatorSerde, because they don’t need protobuf serialization.
Examples: CometBroadcastExchangeExec, CometShuffleExchangeExec
Local TopK inside a sink#
CometTakeOrderedAndProjectExec owns local candidate selection, the optional shuffle, and final
selection with offset and projection. When spark.comet.exec.topK.fusion.enabled is enabled for an
eligible native Parquet scan, its conversion inserts CometLocalTopKExec before native blocks are
serialized. The local node extends CometUnaryExec and serializes a bounded sort above the native
scan, so they execute in one block. It preserves the scan’s output and partitioning and advertises
the local sort order.
The inserted node keeps the original Spark TopK for plan bookkeeping, but only the outer TopK owns
Spark’s offset and projection. Both transition reversion and aggregate-buffer restoration therefore
remove the inserted local node when restoring Spark execution. Restoring its originalPlan would
apply the global TopK twice. A single input partition uses a final native limit and projection;
multiple input partitions still require a final TopK after the shuffle.
With spark.comet.exec.topK.dynamicFilter.enabled, only the local sort sets the protobuf
Sort.dynamic_filter_enabled flag. The native planner wraps an eligible sort in
TopKReaderFilterExec. Its permanent plan contains an unexecuted sort template; each execution
creates a fresh sort and live predicate, then attaches that predicate to the Parquet reader when
eligible. The stream owns the heap and predicate until completion, error, or cancellation. The
final TopK after an exchange does not share this state. Reader work stays on scan metrics while
attachment counters belong to the local TopK.
Choosing the Right Operator Type#
When adding a new operator, choose based on these criteria:
Use Native Operators when:
The operator transforms data (e.g., project, filter, sort, aggregate, join)
The operator has a direct DataFusion equivalent or custom implementation
The operator consumes native child operators and produces native output
The operator is in the middle of an execution pipeline
Use Sink Operators when:
The operator serves as a data source for native execution (becomes a
ScanExec)The operator brings data from non-native sources (e.g.,
UnionExeccombining multiple inputs)The operator is typically a leaf or near-leaf node in the execution tree
The operator needs special handling to interface with the native engine
Implementation Note for Sinks:
Sink operators are handled specially in CometExecRule.operator2Proto. Instead of converting to their own operator
type, they are converted to ScanExec in the native plan. This allows them to serve as entry points for native
execution blocks. The original Spark operator is wrapped with CometScanWrapper or CometSinkPlaceHolder which
manages the boundary between JVM and native execution.
Operators That Should Not Be Converted#
Before adding an operator, check that converting it would speed anything up. Comet deliberately leaves three kinds of Spark plan nodes in place.
Wrappers and scheduling nodes. AdaptiveSparkPlanExec, the AQE query stages
(ShuffleQueryStageExec, BroadcastQueryStageExec, TableCacheQueryStageExec, and, on Spark 4.0
and later, ResultQueryStageExec), AQEShuffleReadExec, InputAdapter, WholeStageCodegenExec,
ReusedExchangeExec, and ReusedSubqueryExec do no data processing of their own. They schedule
stages, mark whole-stage code generation boundaries, choose which shuffle blocks each task reads,
or point at a plan that runs elsewhere. AQE creates the query stages itself, after Comet’s rules
have run on the plan inside them, and depends on their exact class. For example, it casts the root
of the final plan to ResultQueryStageExec. When a query stage wraps a Comet shuffle, broadcast,
or cached relation, CometExecRule reads from it as a native input through CometExchangeSink
and leaves the stage itself in place.
Operators that run user JVM code on JVM objects. The typed Dataset API plans
DeserializeToObjectExec, SerializeFromObjectExec, MapElementsExec, MapPartitionsExec,
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, 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. For the same reason, it leaves the output unconverted when the plan uses
input_file_name(), input_file_block_start() or input_file_block_length(): filling the batch
moves the scan’s reader past the file that these report when Spark evaluates them above the
conversion. 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 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.
If a new Spark version adds a wrapper node, do not write a serde for it. Add it to the nodes that
ExtendedExplainInfo.generateTreeString skips when counting operators, and to the wrapper list in
Understanding Comet Plans, so the coverage
summary does not count it as a Spark operator. If CometExecRule visits the node, also add it to
the operators it leaves in place without recording a fallback reason. ExtendedExplainInfo
already skips every QueryStageExec, so a new query stage type needs no change there.
Implementing a Native Operator#
This section focuses on adding a native operator, which is the most common and complex case.
Step 1: Define the Protobuf Message#
First, add the operator definition to native/proto/src/proto/operator.proto.
Add to the Operator Message#
Add your new operator to the oneof op_struct in the main Operator message:
message Operator {
repeated Operator children = 1;
uint32 plan_id = 2;
oneof op_struct {
Scan scan = 100;
Projection projection = 101;
Filter filter = 102;
// ... existing operators ...
YourNewOperator your_new_operator = 112; // Choose next available number
}
}
Define the Operator Message#
Create a message for your operator with the necessary fields:
message YourNewOperator {
// Fields specific to your operator
repeated spark.spark_expression.Expr expressions = 1;
// Add other configuration fields as needed
}
For reference, see existing operators like Filter (simple), HashAggregate (complex), or Sort (with ordering).
Step 2: Create a CometOperatorSerde Implementation#
Create a new Scala file in spark/src/main/scala/org/apache/spark/sql/comet/ (e.g., CometYourOperatorExec.scala) holding both the CometOperatorSerde[T] object, where T is the Spark operator type, and the CometNativeExec case class it creates. Scan and sink serdes live in spark/src/main/scala/org/apache/comet/serde/operator/ instead.
The CometOperatorSerde trait provides several key methods:
enabledConfig: Option[ConfigEntry[Boolean]]- Configuration to enable/disable this operatorgetSupportLevel(operator: T): SupportLevel- Determines if the operator is supportedconvert(op: T, builder: Operator.Builder, childOp: Operator*): Option[Operator]- Converts to protobufcreateExec(nativeOp: Operator, op: T): CometNativeExec- Creates the Comet execution operator wrapper
The validation workflow in CometExecRule.isOperatorEnabled:
Checks if the operator is enabled via
enabledConfigCalls
getSupportLevel()to determine compatibilityHandles Compatible/Incompatible/Unsupported cases with appropriate fallback messages
Simple Example (Filter)#
import com.google.common.base.Objects
object CometFilterExec extends CometOperatorSerde[FilterExec] {
override def enabledConfig: Option[ConfigEntry[Boolean]] =
Some(CometConf.COMET_EXEC_FILTER_ENABLED)
override def convert(
op: FilterExec,
builder: Operator.Builder,
childOp: OperatorOuterClass.Operator*): Option[OperatorOuterClass.Operator] = {
val cond = exprToProto(op.condition, op.child.output)
if (cond.isDefined && childOp.nonEmpty) {
val filterBuilder = OperatorOuterClass.Filter
.newBuilder()
.setPredicate(cond.get)
Some(builder.setFilter(filterBuilder).build())
} else {
None
}
}
override def createExec(nativeOp: Operator, op: FilterExec): CometNativeExec = {
CometFilterExec(nativeOp, op, op.output, op.condition, op.child, SerializedPlan(None))
}
}
case class CometFilterExec(
override val nativeOp: Operator,
override val originalPlan: SparkPlan,
override val output: Seq[Attribute],
condition: Expression,
child: SparkPlan,
override val serializedPlanOpt: SerializedPlan)
extends CometUnaryExec {
override def outputPartitioning: Partitioning = child.outputPartitioning
override def outputOrdering: Seq[SortOrder] = child.outputOrdering
override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan =
this.copy(child = newChild)
override def stringArgs: Iterator[Any] =
Iterator(output, condition, child)
override def equals(obj: Any): Boolean = {
obj match {
case other: CometFilterExec =>
this.output == other.output &&
this.condition == other.condition && this.child == other.child &&
this.serializedPlanOpt == other.serializedPlanOpt
case _ =>
false
}
}
override def hashCode(): Int = Objects.hashCode(output, condition, child)
}
More Complex Example (Project)#
import com.google.common.base.Objects
object CometProjectExec extends CometOperatorSerde[ProjectExec] {
override def enabledConfig: Option[ConfigEntry[Boolean]] =
Some(CometConf.COMET_EXEC_PROJECT_ENABLED)
override def convert(
op: ProjectExec,
builder: Operator.Builder,
childOp: Operator*): Option[OperatorOuterClass.Operator] = {
val exprs = op.projectList.map(exprToProto(_, op.child.output))
if (exprs.forall(_.isDefined) && childOp.nonEmpty) {
val projectBuilder = OperatorOuterClass.Projection
.newBuilder()
.addAllProjectList(exprs.map(_.get).asJava)
Some(builder.setProjection(projectBuilder).build())
} else {
None
}
}
override def createExec(nativeOp: Operator, op: ProjectExec): CometNativeExec = {
CometProjectExec(nativeOp, op, op.output, op.projectList, op.child, SerializedPlan(None))
}
}
case class CometProjectExec(
override val nativeOp: Operator,
override val originalPlan: SparkPlan,
override val output: Seq[Attribute],
projectList: Seq[NamedExpression],
child: SparkPlan,
override val serializedPlanOpt: SerializedPlan)
extends CometUnaryExec
with PartitioningPreservingUnaryExecNode {
override def producedAttributes: AttributeSet = outputSet
override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan =
this.copy(child = newChild)
override def stringArgs: Iterator[Any] = Iterator(output, projectList, child)
override def equals(obj: Any): Boolean = {
obj match {
case other: CometProjectExec =>
this.output == other.output &&
this.projectList == other.projectList &&
this.child == other.child &&
this.serializedPlanOpt == other.serializedPlanOpt
case _ =>
false
}
}
override def hashCode(): Int = Objects.hashCode(output, projectList, child)
override protected def outputExpressions: Seq[NamedExpression] = projectList
}
Plan Identity and Exchange Reuse#
Spark’s ReuseExchangeAndSubquery identifies equivalent plans through canonicalization and may
reuse one exchange for both branches. If equals omits a parameter that changes results, different
operators can appear equivalent and silently return the wrong rows. If it includes execution state,
equivalent operators can fail to reuse an exchange.
For operators such as Filter and Project, compare the children, output and every parameter that
affects results in equals, and include those semantic fields in hashCode. Capture semantic flags
on the Comet operator itself: a value stored only in the protobuf or original Spark plan cannot
participate in this field-based identity. Keep equals and hashCode consistent: equal operators
must have equal hashes. Do not rely on the default case-class implementations.
The examples above follow the implementations in
operators.scala:
nativeOpis the protobuf representation andoriginalPlanis the source Spark plan; neither is compared by these operators.CometNativeExec.canonicalizePlansclears the non-child Spark plan references, including theiroriginalPlan.serializedPlanOptholds the bytes for a native execution block. These operators compare it inequals, but omit it fromhashCode;CometNativeExec.doCanonicalizeclears the block’s serialized plan. The bytes therefore do not distinguish canonicalized plans.stringArgscontrols the plan’s displayed arguments. Show the semantic fields and children rather than serialization state; this method does not define equality.
Some operators use a different convention. CometNativeScanExec compares originalPlan to retain
scan identity, and its doCanonicalize canonicalizes that plan while removing unused dynamic
pruning filters. CometBroadcastExchangeExec also compares originalPlan before canonicalization,
but its doCanonicalize clears that reference and retains the canonicalized child. Follow each
operator’s equality and canonicalization together rather than copying an exclusion in isolation.
Using getSupportLevel#
Override getSupportLevel to control operator support based on specific conditions:
override def getSupportLevel(operator: YourOperatorExec): SupportLevel = {
// Check for unsupported features
if (operator.hasUnsupportedFeature) {
return Unsupported(Some("Feature X is not supported"))
}
// Check for incompatible behavior
if (operator.hasKnownDifferences) {
return Incompatible(Some("Known differences in edge case Y"))
}
Compatible()
}
Support levels:
Compatible()- Fully compatible with Spark (default)Incompatible()- Supported but may differ; requires explicit opt-inUnsupported()- Not supported under current conditions
Note that Comet will treat an operator as incompatible if any of the child expressions are incompatible.
Step 3: Register the Operator#
Add your operator to the appropriate map in CometExecRule.scala:
For Native Operators#
Add to the nativeExecs map (CometExecRule.scala):
val nativeExecs: Map[Class[_ <: SparkPlan], CometOperatorSerde[_]] =
Map(
classOf[ProjectExec] -> CometProjectExec,
classOf[FilterExec] -> CometFilterExec,
// ... existing operators ...
classOf[YourOperatorExec] -> CometYourOperator,
)
For Sink Operators#
If your operator is a sink (becomes a ScanExec in the native plan), add to the sinks map (CometExecRule.scala):
val sinks: Map[Class[_ <: SparkPlan], CometOperatorSerde[_]] =
Map(
classOf[CoalesceExec] -> CometCoalesceExec,
classOf[UnionExec] -> CometUnionExec,
// ... existing operators ...
classOf[YourSinkOperatorExec] -> CometYourSinkOperator,
)
Note: The allExecs map automatically combines both nativeExecs and sinks, so you only need to add to one of the two maps.
Step 4: Add Configuration Entry#
Add a configuration entry in common/src/main/scala/org/apache/comet/CometConf.scala:
val COMET_EXEC_YOUR_OPERATOR_ENABLED: ConfigEntry[Boolean] =
conf("spark.comet.exec.yourOperator.enabled")
.doc("Whether to enable your operator in Comet")
.booleanConf
.createWithDefault(true)
Run make to update the user guide. The new configuration option will be added to docs/source/user-guide/latest/configs.md.
Step 5: Implement the Native Operator in Rust#
Update the Planner#
In native/core/src/execution/planner.rs, add a match case in the operator deserialization logic to handle your new protobuf message:
use datafusion_comet_proto::spark_operator::operator::OpStruct;
// In the create_plan or similar method:
match op.op_struct.as_ref() {
Some(OpStruct::Scan(scan)) => {
// ... existing cases ...
}
Some(OpStruct::YourNewOperator(your_op)) => {
create_your_operator_exec(your_op, children, session_ctx)
}
// ... other cases ...
}
ProjectionBuilder prunes a DataFusion filter’s output for column-only projections, including
empty projections such as the input to count(*). Each required output column is filtered once,
then the projection restores its order, duplicates and aliases. The predicate still sees its
original input schema. Both native plans remain for metrics, including when they share a Spark
plan ID. Computed projections and projections that use every input column are unchanged.
Implement the Operator#
Create the operator implementation, either in an existing file or a new file in native/core/src/execution/operators/:
use datafusion::physical_plan::{ExecutionPlan, ...};
use datafusion_comet_proto::spark_operator::YourNewOperator;
pub fn create_your_operator_exec(
op: &YourNewOperator,
children: Vec<Arc<dyn ExecutionPlan>>,
session_ctx: &SessionContext,
) -> Result<Arc<dyn ExecutionPlan>, ExecutionError> {
// Deserialize expressions and configuration
// Create and return the execution plan
// Option 1: Use existing DataFusion operator
// Ok(Arc::new(SomeDataFusionExec::try_new(...)?))
// Option 2: Implement custom operator (see ExpandExec for example)
// Ok(Arc::new(YourCustomExec::new(...)))
}
For custom operators, you’ll need to implement the ExecutionPlan trait. Operators that need nothing else from core live in the datafusion-comet-operators crate under native/operators/src/. See native/operators/src/expand.rs or native/core/src/execution/operators/scan.rs for examples.
Step 6: Add Tests#
Scala Integration Tests#
Add tests in spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala or a related test suite:
test("your operator") {
withTable("test_table") {
sql("CREATE TABLE test_table(col1 INT, col2 STRING) USING parquet")
sql("INSERT INTO test_table VALUES (1, 'a'), (2, 'b')")
// Test query that uses your operator
checkSparkAnswerAndOperator(
"SELECT * FROM test_table WHERE col1 > 1"
)
}
}
The checkSparkAnswerAndOperator helper verifies:
Results match Spark’s native execution
Your operator is actually being used (not falling back)
Plan Identity Regression Tests#
If the operator has a result-affecting parameter beyond its children and output, add an
exchange-reuse regression. Build two branches that differ in that parameter, verify that both
execute with Comet, and assert the expected distinct results and that sameResult is false.
Also test equivalent branches with fresh expression IDs or aliases: sameResult should be true,
semanticHash values should match, and the executed plan should contain a reused exchange. Cover AQE both enabled and disabled when applicable.
Use the “aggregate canonicalization preserves result expressions and equivalent reuse” tests in
CometAggregateSuite
as a model. Inspect the executed plans: an optimizer rewrite can make the branches differ for an
unrelated reason, allowing a test to pass even when the intended field is missing from equals.
Choose inputs and query shapes that preserve the parameter difference, and traverse adaptive and
query-stage wrappers when checking the native operator and reused exchange.
Rust Unit Tests#
Add unit tests in your Rust implementation file:
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_your_operator() {
// Test operator creation and execution
}
}
Step 7: Update Documentation#
Add your operator to the supported operators list in docs/source/user-guide/latest/compatibility/operators.md or similar documentation.
Implementing a Sink Operator#
Sink operators are converted to ScanExec in the native plan and serve as entry points for native execution. The implementation is simpler than native operators because sink operators extend the CometSink base class which provides the conversion logic.
Step 1: Create a CometOperatorSerde Implementation#
Create a new Scala file in spark/src/main/scala/org/apache/spark/sql/comet/ (e.g., CometYourSinkOperator.scala):
import org.apache.comet.serde.operator.CometSink
object CometYourSinkOperator extends CometSink[YourSinkExec] {
override def enabledConfig: Option[ConfigEntry[Boolean]] =
Some(CometConf.COMET_EXEC_YOUR_SINK_ENABLED)
// Optional: Override if the data produced is FFI safe
override def isFfiSafe: Boolean = false
override def createExec(
nativeOp: OperatorOuterClass.Operator,
op: YourSinkExec): CometNativeExec = {
CometSinkPlaceHolder(
nativeOp,
op,
CometYourSinkExec(op, op.output, /* other parameters */, op.child))
}
// Optional: Override getSupportLevel if you need custom validation beyond data types
override def getSupportLevel(operator: YourSinkExec): SupportLevel = {
// CometSink base class already checks data types in convert()
// Add any additional validation here
Compatible()
}
}
/**
* Comet implementation of YourSinkExec that supports columnar processing
*/
case class CometYourSinkExec(
override val originalPlan: SparkPlan,
override val output: Seq[Attribute],
/* other parameters */,
child: SparkPlan)
extends CometExec
with UnaryExecNode {
override protected def doExecuteColumnar(): RDD[ColumnarBatch] = {
// Implement columnar execution logic
val rdd = child.executeColumnar()
// Apply your sink operator's logic
rdd
}
override def outputPartitioning: Partitioning = {
// Define output partitioning
}
override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan =
this.copy(child = newChild)
}
Key Points:
Extend
CometSink[T]which provides theconvert()method that transforms the operator toScanExecThe
CometSink.convert()method (inCometSink.scala) automatically handles:Data type validation
Conversion to
ScanExecin the native planSetting FFI safety flags
You must implement
createExec()to wrap the operator appropriatelyYou typically need to create a corresponding
CometYourSinkExecclass that implements columnar execution
Step 2: Register the Sink#
Add your sink to the sinks map in CometExecRule.scala:
val sinks: Map[Class[_ <: SparkPlan], CometOperatorSerde[_]] =
Map(
classOf[CoalesceExec] -> CometCoalesceExec,
classOf[UnionExec] -> CometUnionExec,
classOf[YourSinkExec] -> CometYourSinkOperator,
)
Step 3: Add Configuration#
Add a configuration entry in CometConf.scala:
val COMET_EXEC_YOUR_SINK_ENABLED: ConfigEntry[Boolean] =
conf("spark.comet.exec.yourSink.enabled")
.doc("Whether to enable your sink operator in Comet")
.booleanConf
.createWithDefault(true)
Step 4: Add Tests#
Test that your sink operator correctly feeds data into native execution:
test("your sink operator") {
withTable("test_table") {
sql("CREATE TABLE test_table(col1 INT, col2 STRING) USING parquet")
sql("INSERT INTO test_table VALUES (1, 'a'), (2, 'b')")
// Test query that uses your sink operator followed by native operators
checkSparkAnswerAndOperator(
"SELECT col1 + 1 FROM (/* query that produces YourSinkExec */)"
)
}
}
Important Notes for Sinks:
Sinks extend the
CometSinkbase class, which provides theconvert()method implementationThe
CometSink.convert()method automatically handles conversion toScanExecin the native planYou don’t need to add protobuf definitions for sink operators - they use the standard
ScanmessageYou don’t need Rust implementation for sinks - they become standard
ScanExecoperators that read from the JVMSink implementations should provide a columnar-compatible execution class (e.g.,
CometCoalesceExec)The
createExec()method wraps the operator withCometSinkPlaceHolderto manage the JVM-to-native boundarySee
CometCoalesceExec.scalaorCometUnionExecinspark/src/main/scala/org/apache/spark/sql/comet/for reference implementations
Implementing a JVM Operator#
For operators that run in the JVM:
Create a new operator class extending appropriate Spark base classes in
spark/src/main/scala/org/apache/comet/Add matching logic in
CometExecRule.scalato transform the Spark operatorNo protobuf or Rust implementation needed
Example pattern from CometExecRule.scala:
case s: ShuffleExchangeExec =>
CometShuffleExchangeExec.shuffleSupported(s) match {
case Some(CometNativeShuffle) =>
CometShuffleExchangeExec(s, shuffleType = CometNativeShuffle)
case Some(CometColumnarShuffle) =>
CometShuffleExchangeExec(s, shuffleType = CometColumnarShuffle)
case None => s
}
Common Patterns and Helpers#
Expression Conversion#
Use QueryPlanSerde.exprToProto to convert Spark expressions to protobuf:
val protoExpr = exprToProto(sparkExpr, inputSchema)
Restoring the Spark operator (sparkFallback)#
CometExec.originalPlan is the Spark operator this node replaced. CometExecRule copies
originalPlan.logicalLink onto the Comet node, which is how AQE finds the node again when it
re-plans a stage. RevertNativeForTransitionHeavyStages calls sparkFallback(newChildren) to
rebuild that Spark operator with the children of the reverted stage.
The default implementation is originalPlan.withNewChildren(newChildren). It refuses a null
originalPlan, an originalPlan that is one of the node’s own children, or a different number of
children than the Spark operator has.
Override sparkFallback when conversion changes the plan shape, so the restored node is not that
Spark operator with the same children. CometNativeWriteExec replaces a DataWritingCommandExec
and drops the WriteFilesExec under it; its override puts that wrapper back around the restored
input. CometIcebergWriteExec keeps the same shape as IcebergWriteExec, so the default is
enough.
Also override sparkFallback when the operator’s live state differs from originalPlan.
CometNativeScanExec restores its current partition and data filters, and
CometIcebergNativeScanExec restores its current runtime filters, so AQE’s executable DPP
subqueries survive reversion. CometLocalTopKExec returns the restored child directly: Comet
inserted that local candidate selection, and only the outer TopK restores Spark’s offset and
projection. Rebuilding the original TopK at both nodes would apply it twice.
Do not point originalPlan at a child. If that child is a shuffle or query stage, the copied
logical link puts this node inside the stage’s LogicalQueryStage. AQE then re-plans a second
copy of the operator around the one that is already there.
Handling Fallback#
Use withInfo to tag operators with fallback reasons:
if (!canConvert) {
withInfo(op, "Reason for fallback", childNodes: _*)
return None
}
Child Operator Validation#
Always check that child operators were successfully converted:
if (childOp.isEmpty) {
// Cannot convert if children failed
return None
}
Debugging Tips#
Enable verbose logging: Set
spark.comet.explain.format=verboseto see detailed plan transformationsCheck fallback reasons: Set
spark.comet.explain.fallback.log.enabled=trueto log why operators fall back to SparkVerify protobuf: Add debug prints in Rust to inspect deserialized operators
Use EXPLAIN: Run
EXPLAIN EXTENDEDon queries to see the physical plan