Operator Tuning#

Optimizing Joins#

Spark often chooses SortMergeJoin over ShuffledHashJoin for stability reasons. If the build-side of a ShuffledHashJoin is very large then it could lead to OOM in Spark.

Vectorized query engines tend to perform better with ShuffledHashJoin, so for best performance it is often preferable to configure Comet to convert SortMergeJoin to ShuffledHashJoin. Comet does not yet provide spill-to-disk for ShuffledHashJoin so this could result in OOM. Also, SortMergeJoin may still be faster in some cases. It is best to test with both for your specific workloads.

To configure Comet to convert SortMergeJoin to ShuffledHashJoin, set spark.comet.exec.forceShuffledHashJoin=true.

Join Runtime Filters#

Set spark.comet.exec.join.dynamicFilter.enabled=true to try experimental native hash join runtime filtering. It is disabled by default. Eligible joins are inner joins with one direct signed integer key (TINYINT, SMALLINT, INT, or BIGINT) and one native partition per input within each task. Both broadcast and shuffled hash joins support either Spark build side. Unsupported joins keep their existing execution path.

Once the build completes, its key domain filters probe batches before the hash probe. Eligible native Parquet readers also use the domain to prune row groups. Reader attachment can pass through direct-column IS NOT NULL checks, including conjunctions, and remaps columns when the scan itself projects the file schema. The original null checks and residual runtime filter remain in place. The original join still verifies matches, including any hash collisions admitted by the filter. Standalone projections, other filter expressions, and limits prevent reader attachment.

To preserve schema-conversion and timestamp-overflow errors, runtime reader pruning is disabled for each file whose projected or statically filtered columns require schema adaptations beyond direct column mappings or literal values. This conservative check also disables reader pruning for allowed INT32 to BIGINT promotion and for projecting a subset of a struct’s fields, even when those adaptations cannot fail. Nested column pruning still reads only the requested struct fields. Scans with supplied file statistics also skip reader attachment. These cases still use runtime filtering on decoded batches.

Filters stay within the task’s native plan and do not propagate across Spark exchanges or JVM/Arrow boundaries. A shuffled hash join can still filter probe batches after shuffle, but it cannot send its filter back to an earlier scan stage. Compare the runtime-filter and scan metrics with the setting disabled to distinguish reduced hash-probe work from reader I/O savings.

Adaptive Partial Aggregation#

Set spark.comet.exec.aggregate.skipPartial.enabled=true to let Comet bypass partial hash aggregation for high-cardinality grouping when it is not reducing the number of rows enough. This experimental optimization is disabled by default. It currently applies only to fused native shuffle-writer plans whose partial aggregates are grouping-only or single-argument COUNT. Low-cardinality inputs continue to aggregate normally. The SQL metric rows bypassing partial aggregation shows whether skipping occurred.

DataFusion makes the decision separately in each task. It starts checking after the first 100,000 input rows, and as soon as the number of groups divided by the number of input rows exceeds 0.8, it stops aggregating and sends the rest of the task’s rows to the shuffle as they are. It does not check again, so a task whose keys repeat after a mostly distinct start, such as several snapshot files of the same keys packed into one split, can shuffle many times more rows than it would with skipping disabled. Compare the shuffle write metrics with the setting enabled and disabled before enabling it for a workload.

Eligibility is conservative for the whole fused native plan: any unsupported partial accumulator, Spark PartialMerge, or mixed-mode aggregate disables skipping in that plan. Multi-argument COUNT and other accumulators are not admitted. Distribution-required grouping-only stages still fully deduplicate, and non-native-shuffle plans retain ordinary aggregation.

To experiment with the thresholds, also enable spark.comet.exec.respectDataFusionConfigs, a development and testing option that defaults to false. For example, the following SQL settings pass through the default threshold values, which you can adjust:

SET spark.comet.exec.aggregate.skipPartial.enabled=true;
SET spark.comet.exec.respectDataFusionConfigs=true;
SET spark.comet.datafusion.execution.skip_partial_aggregation_probe_rows_threshold=100000;
SET spark.comet.datafusion.execution.skip_partial_aggregation_probe_ratio_threshold=0.8;

A lower row threshold allows an earlier decision; a lower ratio threshold makes skipping more likely. These settings only tune eligible plans. They cannot enable skipping while spark.comet.exec.aggregate.skipPartial.enabled is false, or for unsupported accumulators and modes.

Local TopK Fusion#

Set spark.comet.exec.topK.fusion.enabled=true to run an eligible local TopK in the same native execution as its Parquet scan. This experimental optimization is disabled by default. It currently supports a direct native Parquet scan ordered by one signed integer column (TINYINT, SMALLINT, INT, or BIGINT). Both sort directions and null orderings are supported. Other inputs use the existing TopK execution path.

Each scan partition keeps enough candidates for both LIMIT and OFFSET. With multiple partitions, Comet shuffles those candidates and performs the final TopK. With one partition, it reuses the local ordering without building a second heap. The final stage applies the offset and output projection.

Fusion reduces the work of passing scan batches between native execution blocks. Fusion alone still reads all input rows. TopK reader pruning is enabled separately. Fusion can also reduce overlap between scan decoding and TopK processing, so some workloads may run slower. Compare enabled and disabled runs with your data layout, payload width, limit, and partition count before enabling it. The CometTopKBenchmark microbenchmark covers these cases with ascending, descending, and random layouts.

TopK Reader Pruning#

Set both spark.comet.exec.topK.fusion.enabled=true and spark.comet.exec.topK.dynamicFilter.enabled=true to pass the local TopK’s improving threshold to its Parquet reader. Both options are experimental and disabled by default. Eligibility is the same single signed integer key described above. Each task creates a fresh threshold; it is not shared across Spark partitions or exchanges, or retained for later executions.

Once the heap contains enough candidates for LIMIT + OFFSET, the reader can skip later row groups whose statistics prove that no row can improve those candidates. Existing Parquet page-index and decoder-filter options can also use the predicate. This option adds no separate filter over decoded scan batches. TopK continues to select the final candidates.

Reader attachment is conservative. A scan with a fetch limit, supplied file statistics, or a static predicate other than direct column IS NOT NULL checks keeps the existing execution path. For each file, schema adaptation disables pruning if it could hide a conversion error in a projected or filtered column. Missing null counts remain unknown, which can prevent pruning even when min/max statistics are present. These cases can still execute a fused TopK.

Reader pruning is most useful when small K values and the file order establish a strong threshold early. Descending or random layouts for an ascending query can prune few or no groups, while still paying the cost of attaching and checking the predicate. Wider rows can increase the benefit when groups are skipped. Compare pruning with fused in CometTopKBenchmark to measure the reader effect, and compare both with unfused to include the cost of fusion. Check the scan’s emitted rows, bytes_scanned, row_groups_pruned_dynamic_filter, and row_groups_pruned_statistics alongside elapsed time; attachment alone does not demonstrate a saving. Pruning when later files open uses the TopK threshold already available and increments row_groups_pruned_statistics. With one row group per file, the dynamic counter can stay zero despite substantial TopK pruning. The statistics counter also includes other predicates, so compare with filtering disabled to assess TopK savings. See TopK metrics.

Optimizing Sorting on Floating-Point Values#

Comet normalizes NaN payloads and signed zeros in scalar FLOAT and DOUBLE ordering keys, so ORDER BY, window ordering and range partitioning on them match Spark and stay native even with spark.comet.exec.strictFloatingPoint=true. Only the comparison key is normalized; returned values keep their original NaN representation and zero sign.

Floating-point values nested in arrays, structs, or maps are compared with Arrow’s raw total ordering instead, which can differ from Spark when the data contains both zero and negative zero, or more than one NaN representation. This is likely an edge case that is not of concern for many users. Setting spark.comet.exec.strictFloatingPoint=true makes those nested cases fall back to Spark, and they can be forced back onto the native path with spark.comet.expression.SortOrder.allowIncompatible=true.

sort_array is separate. It sorts array elements rather than ordering rows, and its elements are compared with Arrow’s raw total ordering, so spark.comet.exec.strictFloatingPoint=true makes it fall back even for a scalar floating-point element type. Use spark.comet.expression.SortArray.allowIncompatible=true to keep it native.