Repository navigation
chore: joom-1.1-7 batch of #10 #11 #12 #13 #14 - #15
Merged
Merged
Conversation
… codegen A Spark BroadcastHashJoinExec can read a Comet broadcast through CometColumnarToRowExec, whose doExecuteBroadcast turns the Arrow batches into the join's relation. This happens under AQE when a re-optimization leaves the join in Spark after its build side was already planned as a Comet broadcast stage, for example when OptimizeSkewedJoin puts an AQEShuffleRead under the sort-merge join below it. When the join itself is outside whole-stage codegen (more output fields than spark.sql.codegen.maxFields), CollapseCodegenStages wrapped the CometColumnarToRowExec into a WholeStageCodegenExec of its own, and the join failed with "WholeStageCodegen (n) does not implement doExecuteBroadcast". CometColumnarToRowExec over a Comet broadcast exchange, a broadcast query stage or a reused exchange of one now reports supportCodegen = false, so it stays the join's build plan, or sits behind an InputAdapter that forwards executeBroadcast when the join is codegen'd. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…receding conjuncts DataFusion's AND evaluates its right side on every row once more than 20% of rows pass the left side. A conjunct that runs through the JVM codegen dispatcher (ScalaUDF, regex family outside the native subset, other dispatched expressions) then costs a JNI round trip and Arrow export for rows Spark would never evaluate it on. CometFilterExec now splits the serialized deterministic predicate into its top-level conjuncts and rewrites `a AND b`, where b contains a JvmScalarUdf, into `CASE WHEN a THEN b END`. DataFusion's CaseExpr projects the columns b needs, filters them by a and evaluates b on the selected rows only. Null results reject rows just like false, so the filter output is unchanged. Scan pushdown is untouched because only the native filter predicate changes. Controlled by spark.comet.exec.filter.shortCircuitJvmDispatch.enabled (default true). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…aggregate Spark plans MIN or MAX over a string as a SortAggregate, since the aggregation buffer is not mutable, and Comet did not convert SortAggregate, so such an aggregate and often the joins around it ran in Spark. MIN and MAX now accept StringType with the default UTF8_BINARY collation. The native min/max compare UTF-8 bytes, which is Spark's binary order; collated strings stay in Spark. A SortAggregateExec is converted to a CometHashAggregateExec built from an ObjectHashAggregateExec with the same fields. Spark runs that operator over any input order, so an aggregate reverted to Spark later (cost-based engine choice, unsafe partial aggregates) stays correct, and the cost-based choice prices it as an object hash aggregate. A sort by exactly the grouping keys directly below the aggregate is dropped; without one only order-insensitive aggregates (MIN, MAX, COUNT, SUM, AVG) are converted. A native sort above the aggregate restores the output ordering a consumer such as a sort-merge join may rely on, and is dropped again when a shuffle reads it directly. Aggregates Comet does not support keep the SortAggregate. spark.comet.exec.sortAggregate.enabled (default true) turns the conversion off. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…utput spark.comet.exec.sort.spillBeforeOutputThreshold (568955d) protects a Spark sort, sort-merge join or window that reads a native sort's output through a columnar-to-row transition in the same task: the native sort cannot be made to release memory, so it spills its in-memory input before producing output and holds only the merge buffers. The threshold was a session extension, so every ExternalSorter in the native plan applied it, including sorts whose output never leaves the plan: under a window, a window group limit, a sort-merge join, an aggregate or the native shuffle writer. Those sorts wrote and read back their whole input once per task for nothing, because their consumer reserves memory from the same native pool and the sort releases it as its output is merged. SortExec now carries a spill_before_output flag, kept by with_fetch, projection swaps and filter pushdown, and ExternalSorter applies the threshold only when it is set. The planner sets it for a Sort reached from the root of the native plan through streaming operators only (Projection, Filter, Limit, Sample, Expand, Explode), whose output is what the JVM consumes. Window, partition-aggregate and sorted window operators do not use ExternalSorter and are unaffected. On a 1/40 sample of sat_product_info_loader's main stage (scan, window, sort-merge join, union, sort, native shuffle) with the joom-1.1-6 image and the patched library loaded from java.library.path, the query took 0.736 task-h with no spill, the same as threshold 0 (0.733-0.756), against 0.824-0.839 task-h and 29.2 GB spilled before. On the 2026-10-10 night, CometSort operators under native consumers spilled 21.6 TB, about one spill per task, mostly under window group limit, window, sort-merge join and native shuffle. Tests: the planner marks only sorts reached through streaming operators, and a planned SortExec carries the flag; the native sort with a JVM consumer still spills before output at the threshold, and with a native consumer it keeps its input in memory; in CometExecSuite, with a 1 byte threshold, sorts read through columnar-to-row spill and sorts under a window and a sort-merge join do not, all matching Spark. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
… needs it The native sort above a SortAggregate converted to a hash aggregate was added unconditionally and dropped only under a shuffle. It is now added the way EnsureRequirements would: when a consumer is visited, before it is converted, each child whose required ordering is no longer met and that leads through unary operators to a converted SortAggregate gets a sort on that edge. A partial aggregate under a shuffle, a final aggregate read by a hash aggregate, a project or a write, and any consumer without an ordering requirement get none. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Spark's sort is stable, so a SortAggregate over input sorted by its grouping keys sees each group's rows in their input order, and FIRST, LAST, ANY_VALUE and FIRST_VALUE return the first row of a group as read. The native hash aggregate does not guarantee that order, so converting such a SortAggregate broke the exact first/last results the fork promises on an ordered single split (CometFirstLastBoolAggFuzzSuite, 240 mismatches on Spark 4.1 and 559 on 3.5 for first, any_value and first_value over strings). A SortAggregate is now converted only when every aggregate function is order-insensitive (MIN, MAX, COUNT, SUM, AVG), whatever its input. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
MAX_BY and MIN_BY accept a StringType value with the default UTF8_BINARY collation; the ordering stays limited to fixed-length types. The native MaxMinBy accumulators already store the value as Arrow row bytes and only compare the ordering, so no native change is needed. A string value makes Spark plan a SortAggregate, which is now converted when its other aggregates are order-insensitive. MAX_BY and MIN_BY depend on the input order only among rows tied on the ordering, where Spark is non-deterministic too and the native path already documents that it may pick a different row. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…t' into joom/1.1-7-rc
…umer' into joom/1.1-7-rc
msafonov
marked this pull request as ready for review
October 10, 2026 20:46
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Release candidate joom-1.1-7-rc1 (tag
joom-1.1-7-rc1):branch-1.1with PRs merged as one batch.Draft: opened to run CI on the combination. Real-data compares and the full benchmark suite run on image
msaf-spark-comet:joom-1.1-7-rc1. Do not merge until approved.🤖 Generated with Claude Code