Repository navigation
docs: prefer behavior-named parameters over Spark versions #6776
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
| native code carries no version logic (#6740). Flag a version-named parameter, and comments that say | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Should the rule also cover native function names? |
||
| "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 | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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. | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
|
||
|
|
||
| #### 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. | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 |
||
| - Spark occasionally changes behavior in a patch release or between minor releases, so a version number is a poor proxy for the behavior. | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The |
||
| - 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. | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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. | ||
|
|
||
There was a problem hiding this comment.
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 sharedserde/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.