Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion .ai/skills/implement-comet-expression/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,8 @@ Follow `adding_a_new_expression.md`:
2. Register it in the matching map in `QueryPlanSerde.scala`.
3. If the function name collides with a DataFusion built-in that has a different signature, use `scalarFunctionExprToProtoWithReturnType` (see "When to set the return type explicitly").
4. For a new scalar function, add a match case in `native/spark-expr/src/comet_scalar_funcs.rs::create_comet_physical_fun`. If step 2 found an upstream implementation, wire that in. Otherwise implement the function under `native/spark-expr/src/`.
5. Add at least one Comet SQL Test at `spark/src/test/resources/sql-tests/expressions/<category>/$ARGUMENTS.sql` exercising column references, literals, and `NULL`.
5. If Spark's behavior differs between versions, resolve the version in the serde or a shim and pass native code a parameter named for the behavior (`wrap_second_millisecond_overflow`), not the version (`spark_420_plus`). See "Name the behavior, not the Spark version" in `adding_a_new_expression.md`.
6. Add at least one Comet SQL Test at `spark/src/test/resources/sql-tests/expressions/<category>/$ARGUMENTS.sql` exercising column references, literals, and `NULL`.

Build and smoke-test:

Expand Down
4 changes: 3 additions & 1 deletion .ai/skills/review-comet-expression-pr/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -266,7 +266,9 @@ reference in that doc is the only place they are documented.
3. **Wrong return type**, it must match Spark exactly
4. **Tests in the wrong framework**, Scala tests where a SQL file test would do
5. **Missing `getSupportLevel`**, divergences left undeclared rather than marked `Incompatible`
6. **Version-specific Spark behavior implemented once**, with no shim
6. **Version-specific Spark behavior implemented once**, with no shim, or a shim whose protobuf
field or native parameter is named for the Spark version (`spark_420_plus`) instead of the
behavior (`wrap_second_millisecond_overflow`)
7. **Name collides with a DataFusion built-in** and no explicit return type
8. **Timestamp result mislabelled**, with the session timezone or no timezone instead of `"UTC"`.
It only shows once the result is compared or fed to another expression.
6 changes: 6 additions & 0 deletions .ai/skills/review-comet-pr/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,12 @@ Comet supports several Spark versions. Version-specific behavior belongs in the
version string in shared code, and not in native Rust. If the PR adds a shim for one 4.x version,
check that the sibling 4.x source sets got it too.

Parameters, protobuf fields, and native function arguments that vary with the Spark version should
describe the behavior, such as `wrap_second_millisecond_overflow`, not the version, such as
`spark_420_plus`. Forks that backport fixes can then set the flag from their own shim, and the
Comment on lines +179 to +181

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.

The paragraph just above says version-specific behavior belongs in the shims and "not in branches on a version string in shared code". The guide added here says to resolve the version "in the serde or a shim", and its example is builder.setWrapSecondMillisecondOverflow(!isSpark42Plus) in the shared serde/datetime.scala. A reviewer following both paragraphs would flag #6740's own pattern.

Could we reconcile the two? One way is to say that the version check may live in Scala, in a shim or as a named predicate in the serde like RegrSparkVersions.slopeFiltersVarByPairNulls, and that what this rule protects is that proto and native code only ever see the behavior.

native code carries no version logic (#6740). Flag a version-named parameter, and comments that say

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.

Cast.is_spark4_plus in expr.proto already breaks this rule, and its comment says it gates more than one behavior ("such as the handling of leading whitespace before T-prefixed time-only strings"). Once this lands, any PR that touches Cast will get flagged for it. Could we file a tracking issue to split it into behavior-named fields and link it here? Renaming the field keeps its number, so it stays wire-compatible because the JVM and native sides ship together.

Should the rule also cover native function names? get_json_object_spark34 is picked by the Spark 3.4 shim and has the same problem in a function name.

"Spark 4.2" for behavior later releases inherit.

When Spark changed the behavior in a patch release, such as SPARK-55969 or SPARK-54918, a check on
the minor version is wrong for every earlier patch. CI builds only the newest patch of each line, so
it can't catch that (#6042, #5701). The pull request CI also runs only the default Spark profile, so
Expand Down
2 changes: 2 additions & 0 deletions .ai/skills/wire-datafusion-function/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,8 @@ Helpers from `QueryPlanSerde`:
- `optExprWithInfo(optExpr, expr, children*)` — wrap final result; propagates "why we couldn't convert" tags.
- `withInfo(expr, "reason")` — tag a fallback when returning `None`.

If the Spark behavior you are matching varies by Spark version, resolve the version in Scala and pass the native side a parameter named for the behavior, not the version (see "Name the behavior, not the Spark version" in `adding_a_new_expression.md`).

`getSupportLevel` returning `Incompatible(Some("…"))` gates behind `spark.comet.expr.<name>.allowIncompatible=true`. `Unsupported(…)` always falls back.

### 4. Register the UDF (Pattern B only)
Expand Down
12 changes: 12 additions & 0 deletions docs/source/contributor-guide/adding_a_new_expression.md
Original file line number Diff line number Diff line change
Expand Up @@ -583,6 +583,18 @@ If the expression you're adding has different behavior across different Spark ve
1. Shims that exist in `spark/src/main/spark-$SPARK_VERSION/org/apache/comet/shims/CometExprShim.scala` for each Spark version. These shims are used to provide compatibility between different Spark versions.
2. Variables that correspond to the Spark version, such as `isSpark33Plus`, which can be used to conditionally execute code based on the Spark version.

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.

isSpark33Plus no longer exists. The helpers today are isSpark35Plus through isSpark42Plus. Since the new section builds directly on this list, could we update it to isSpark42Plus here?


#### Name the behavior, not the Spark version

When a code path or a protobuf field depends on how Spark behaves, name it for the behavior and not for the Spark version that introduced it. Resolve the version in Scala, in the serde or a shim, and pass the result to native code as a parameter that describes the behavior. For example, `TruncTimestamp` carries a `wrap_second_millisecond_overflow` flag, which the serde sets from `!isSpark42Plus`. It is not called `spark_420_plus`.

Reasons to prefer behavior names:

- Some deployments run a Spark fork that backports fixes from newer open-source releases. They can set a behavior flag in their own shim, but cannot make a `spark_420_plus` flag mean something it does not.

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.

A fork can only set the flag "in their own shim" if the flag is resolved in a shim. With #6740 it is resolved in the shared datetime.scala, so a fork still patches shared Scala code. It just doesn't have to touch the proto or the Rust. Could we state the benefit as what it is today, or recommend resolving the flag in a shim or a single named predicate so that a fork overrides one place?

- Spark occasionally changes behavior in a patch release or between minor releases, so a version number is a poor proxy for the behavior.

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.

The TruncTimestamp example resolves at minor-version granularity, so it doesn't show this case. RegrSparkVersions.r2DegenerateCasesSwapped in serde/aggregates.scala resolves a behavior flag down to the patch version for SPARK-55969 (3.5.9, 4.0.3, 4.1.2, 4.2.0). Would it be worth pointing to it as the model when a fix lands in patch releases? filter_var_by_pair_nulls is also a good existing example of a behavior-named proto field.

- The native code stays free of Spark version logic, and a reader of the Rust or proto definition can tell what the flag does without looking up a release.

Write comments the same way. Say "Spark 4.2 and later" for a behavior that future releases inherit, rather than "Spark 4.2". If a later release changes the behavior again, add a new parameter or shim for that change.

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.

Should we also ask comments to cite the SPARK JIRA, as #6740 does with SPARK-56663? Forks track backports by JIRA rather than by release, so it supports the fork rationale above more directly than the version does.


## Shimming to Support Different Spark Versions

If the expression you're adding has different behavior across different Spark versions, you can use the shim system located in `spark/src/main/spark-$SPARK_VERSION/org/apache/comet/shims/CometExprShim.scala` for each Spark version.
Expand Down