What / Why
We spend a lot of effort profiling and optimizing native execution, but comparatively little on the JVM-side planner: CometScanRule / CometExecRule, the QueryPlanSerde expression and operator serde, and the post-columnar rules. That cost is paid on the driver for every query, and again for every query stage under AQE (see the rule injection points documented in CometSparkSessionExtensions), so it shows up as latency on short queries and on workloads with many small stages.
This EPIC collects findings from a first pass over that code. Measurements below are from a local probe on a modest 2-table join + aggregate query (778-char SQL, 90 Expr nodes in the serialized plan), Spark 4.1 / JDK 17, debug native build:
executedPlan w/ Comet = 23.53 ms
executedPlan w/o Comet = 19.05 ms
CometScanRule = 0.23 ms
CometExecRule = 1.13 ms
total serialized plan bytes = 58445
query_context serialized bytes = 56227 (96.2% of plan)
So the rules themselves are a modest ~4.5 ms of overhead here, but the serialized plan is 32x larger than it needs to be, which is paid per task on every executor.
Findings
1. Serialized plan duplicates the full SQL text on every expression
Tracked separately in #5200 — 96% of plan bytes, 32x plan-size inflation, 2.4x per-task protobuf round-trip cost. This is the largest item by a wide margin.
2. Uncached Class.getMethod throughout the Iceberg reflection paths
Likely the largest planner-time cost on Iceberg workloads. IcebergReflection.extractFileLocation resolves location() on every file and detects older Iceberg by catching NoSuchMethodException, so on path()-only versions every call constructs an exception with a stack trace. It is called once per file scan task (CometIcebergNativeScan.scala, in the tasksCollection.asScala.foreach { task => ... } loop) and once per delete file. In the same loop, extractDeleteFilesList re-resolves keyMetadata() per task, and serializePartitionInfo resolves spec, partitionType, fields, fieldId and get per task and per partition field. CometScanRule.validateIcebergFileScanTasks does getMethod("transform") + setAccessible inside its per-task loop. IcebergReflection.getFieldInfo resolves three methods per column.
All of these resolve to a fixed handful of Method objects. Class.getMethod does a linear search and allocates a fresh Method copy on each call, so a scan with tens of thousands of file tasks does O(100k) reflective lookups per planning pass — repeated per AQE stage. There is already a // TODO acknowledging this in IcebergReflection.getFileFormat. Hoisting the lookups into lazy vals / passing already-resolved Methods down is mechanical.
3. tagUnsafePartialAggregates serializes aggregates only to throw the result away
CometExecRule.canAggregateBeConverted calls QueryPlanSerde.exprToProto / aggExprToProto on grouping, aggregate and result expressions purely to get a Boolean. The protos are discarded and then rebuilt during the real transform pass, and the probe also burns exprIdCounter entries and tags fallback reasons speculatively.
The guards short-circuit, so this only fires for Final aggregates whose functions are not mixed-execution-safe, and for multi-stage collect_list / collect_set — it is not on the path of every query. But when it does fire it is a full double serialization, once per AQE stage. The method's own docstring already notes that a shared predicate helper would be preferable.
4. DecimalPrecision.promote is a second full traversal of every expression tree
QueryPlanSerde.exprToProto runs expr.transformUp over the whole tree before the serde walks it again. A cheap guard (does the tree contain a decimal BinaryArithmetic at all?) or folding the promotion into the serde walk would remove one of the two traversals. Small per call, but it is called once per projection / filter / aggregate expression.
5. EliminateRedundantTransitions.hasCometNativeChild is quadratic
It does a full op.exists(...) subtree scan at each ColumnarToRowExec encountered during a transformUp, so cost is quadratic in the number of C2R nodes. Bottom-up memoization is straightforward.
6. PlanDataInjector.injectPlanData rebuilds the whole operator tree per task
It recurses into every subtree and rebuilds every operator via builders even when only one leaf scan is actually injected. Returning the original op unchanged when no descendant needs injection would leave only the path to the scan to rebuild. This runs per task on the executor.
7. Micro-optimizations
CometConf.getBooleanConf builds a config-key string, does a map lookup and allocates a toLowerCase copy — twice per expression (isExprEnabled + isExprAllowIncompat).
ExpressionRegistry::can_handle followed by create_expr computes get_expression_type and hashes twice where a single get would do (native/core/src/execution/planner/expression_registry.rs).
CometScanRule allocates the Iceberg metadata-table name Set on every rule invocation.
Not investigated yet
- Whether the Comet rules can cheaply detect "already converted" on AQE stage re-entry and skip more work.
- Cost of
CometScanRule on wide schemas (isSchemaSupported walks the full schema per scan, per rule invocation).
- Whether a committed planner benchmark belongs in the repo so these are tracked over time. The probe used for the numbers above was throwaway; turning it into a proper benchmark is worth doing.
Reproducing
The numbers above came from a temporary ScalaTest suite that plans a query repeatedly with the Comet configs toggled (interleaving the configs across rounds so JIT warm-up does not bias whichever config is measured first — this matters a lot here), reads CometNativeExec.serializedPlanOpt, and walks the parsed Operator proto counting Expr nodes and query_context bytes.
What / Why
We spend a lot of effort profiling and optimizing native execution, but comparatively little on the JVM-side planner:
CometScanRule/CometExecRule, theQueryPlanSerdeexpression and operator serde, and the post-columnar rules. That cost is paid on the driver for every query, and again for every query stage under AQE (see the rule injection points documented inCometSparkSessionExtensions), so it shows up as latency on short queries and on workloads with many small stages.This EPIC collects findings from a first pass over that code. Measurements below are from a local probe on a modest 2-table join + aggregate query (778-char SQL, 90
Exprnodes in the serialized plan), Spark 4.1 / JDK 17, debug native build:So the rules themselves are a modest ~4.5 ms of overhead here, but the serialized plan is 32x larger than it needs to be, which is paid per task on every executor.
Findings
1. Serialized plan duplicates the full SQL text on every expression
Tracked separately in #5200 — 96% of plan bytes, 32x plan-size inflation, 2.4x per-task protobuf round-trip cost. This is the largest item by a wide margin.
2. Uncached
Class.getMethodthroughout the Iceberg reflection pathsLikely the largest planner-time cost on Iceberg workloads.
IcebergReflection.extractFileLocationresolveslocation()on every file and detects older Iceberg by catchingNoSuchMethodException, so onpath()-only versions every call constructs an exception with a stack trace. It is called once per file scan task (CometIcebergNativeScan.scala, in thetasksCollection.asScala.foreach { task => ... }loop) and once per delete file. In the same loop,extractDeleteFilesListre-resolveskeyMetadata()per task, andserializePartitionInforesolvesspec,partitionType,fields,fieldIdandgetper task and per partition field.CometScanRule.validateIcebergFileScanTasksdoesgetMethod("transform")+setAccessibleinside its per-task loop.IcebergReflection.getFieldInforesolves three methods per column.All of these resolve to a fixed handful of
Methodobjects.Class.getMethoddoes a linear search and allocates a freshMethodcopy on each call, so a scan with tens of thousands of file tasks does O(100k) reflective lookups per planning pass — repeated per AQE stage. There is already a// TODOacknowledging this inIcebergReflection.getFileFormat. Hoisting the lookups into lazy vals / passing already-resolvedMethods down is mechanical.3.
tagUnsafePartialAggregatesserializes aggregates only to throw the result awayCometExecRule.canAggregateBeConvertedcallsQueryPlanSerde.exprToProto/aggExprToProtoon grouping, aggregate and result expressions purely to get aBoolean. The protos are discarded and then rebuilt during the realtransformpass, and the probe also burnsexprIdCounterentries and tags fallback reasons speculatively.The guards short-circuit, so this only fires for
Finalaggregates whose functions are not mixed-execution-safe, and for multi-stagecollect_list/collect_set— it is not on the path of every query. But when it does fire it is a full double serialization, once per AQE stage. The method's own docstring already notes that a shared predicate helper would be preferable.4.
DecimalPrecision.promoteis a second full traversal of every expression treeQueryPlanSerde.exprToProtorunsexpr.transformUpover the whole tree before the serde walks it again. A cheap guard (does the tree contain a decimalBinaryArithmeticat all?) or folding the promotion into the serde walk would remove one of the two traversals. Small per call, but it is called once per projection / filter / aggregate expression.5.
EliminateRedundantTransitions.hasCometNativeChildis quadraticIt does a full
op.exists(...)subtree scan at eachColumnarToRowExecencountered during atransformUp, so cost is quadratic in the number of C2R nodes. Bottom-up memoization is straightforward.6.
PlanDataInjector.injectPlanDatarebuilds the whole operator tree per taskIt recurses into every subtree and rebuilds every operator via builders even when only one leaf scan is actually injected. Returning the original
opunchanged when no descendant needs injection would leave only the path to the scan to rebuild. This runs per task on the executor.7. Micro-optimizations
CometConf.getBooleanConfbuilds a config-key string, does a map lookup and allocates atoLowerCasecopy — twice per expression (isExprEnabled+isExprAllowIncompat).ExpressionRegistry::can_handlefollowed bycreate_exprcomputesget_expression_typeand hashes twice where a singlegetwould do (native/core/src/execution/planner/expression_registry.rs).CometScanRuleallocates the Iceberg metadata-table nameSeton every rule invocation.Not investigated yet
CometScanRuleon wide schemas (isSchemaSupportedwalks the full schema per scan, per rule invocation).Reproducing
The numbers above came from a temporary ScalaTest suite that plans a query repeatedly with the Comet configs toggled (interleaving the configs across rounds so JIT warm-up does not bias whichever config is measured first — this matters a lot here), reads
CometNativeExec.serializedPlanOpt, and walks the parsedOperatorproto countingExprnodes andquery_contextbytes.