Repository navigation
Conversation
Spark 4.2.0 (SPARK-56663) raises `long overflow` when SECOND or MILLISECOND truncation falls below the smallest timestamp, where earlier versions wrap the Long subtraction. The serde now tells the native kernel which behavior to follow, and the wrapping fixture from apache#5956 is split into a 4.1-and-earlier file and a 4.2 file.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Native
date_truncwrapped SECOND/MILLISECOND underflow on every Spark version, whereas Spark 4.2 raiseslong overflow. - Design approach: A serialized behavior flag selects wrapping or checked subtraction without introducing Spark-version logic into Rust.
- Correctness / compatibility analysis: Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 confirm the intended version split. All 26 locally selected native tests passed, including overflow boundaries, poisoned NULLs and unused dictionary entries. One reproducible P2 performance regression remains below.
- Key design decisions: The flag participates in expression equality, hashing and reconstruction. Reusing the existing boundary kernel keeps the change contained, but the dictionary decision unnecessarily routes safe non-UTC fine units through key masking.
- Implementation sketch: Scala sets the protobuf flag, the native builder forwards it, and the kernel applies checked subtraction at the lower boundary. Paired SQL fixtures cover the two version ranges.
- Behavioral changes worth calling out: Spark 4.2 now raises instead of wrapping. Comparison with
branch-1.1also identified already-landed changes to extreme-timestamp handling, UTC labels, timezone normalization and dictionary masking. Reviewed all 14 files in the supplied endpoint diff. The four release-documentation reversals come from base-only commitf7eb8aa3, not this PR's commit. - Suggested improvements: Preserve values-only processing for safe SECOND/MILLISECOND dictionaries regardless of timezone, while retaining masking when overflow is possible.
Reviewed full SHA da1d1eedbed5db835148223c651b00fa79bfe821 against f7eb8aa3abd6d74ae74d0fa0d50d98f4ce65002c. Routed skills: review-comet-pr and review-comet-expression-pr, including the timezone guidance. Existing reviews, comments and threads were empty.
Exact-head CI at the final check: 31 successful, 17 running and 51 skipped checks, with no failures. Native build, Rust tests and lint passed. Expression jobs for all five Spark profiles remained running. Spark's own SQL suites were skipped.
Validation limits: Local validation ran cargo test --manifest-path native/Cargo.toml -p datafusion-comet-spark-expr --lib timestamp_trunc --offline --no-default-features --locked. JVM SQL suites and production benchmarks were not run locally. The performance reproduction compiled the exact base/head kernels with rustc -O against identical debug-built dependencies, so its timings describe that bounded kernel comparison, not end-to-end query throughput.
| let granularity = normalize_timestamp_trunc_format(format)?; | ||
| if matches!(granularity, "microsecond" | "millisecond" | "second") { | ||
| if granularity == "microsecond" | ||
| || (wrap_second_millisecond_overflow && matches!(granularity, "millisecond" | "second")) |
There was a problem hiding this comment.
[P2] [P2] Preserve values-only truncation for safe non-UTC dictionaries
On Spark 4.2, with spark.sql.session.timeZone=Asia/Tokyo and spark.comet.expression.TruncTimestamp.allowIncompatible=true, this condition sends SECOND/MILLISECOND dictionaries into the unconditional non-UTC masking branch below. Even ordinary, non-overflowing timestamps now allocate a mask and scan every key. These units are timezone-independent, so safe inputs should still require work proportional to the distinct values. A base/head kernel comparison with 32 distinct timestamps measured approximately 4.9 µs versus 30.8 µs for 8,192 keys, and 4.8 µs versus 211.5 µs for 65,536 keys, with identical results. Could the checked fine-unit path perform its value-range check before the non-UTC guard, retaining key masking only when overflow is possible?
Evidence: Reproduction: /tmp/review-6740/kernel_bench.rs compiles the exact supplied base and head temporal.rs modules together using rustc -O and identical dependency artifacts. Construct DictionaryArray<Int32Type> with keys i % 32, values 1704067200123456 + j * 1234567 for j=0..31, and timezone Asia/Tokyo. Compare base timestamp_trunc_dyn(input, unit) with head timestamp_trunc_dyn(input, unit, false) for SECOND and MILLISECOND. Results were asserted equal. Two alternating rounds of 2,000 calls each consistently showed about 6x higher kernel time at 8,192 keys and 44x at 65,536 keys. UTC controls remained approximately unchanged. Timings use debug-built dependencies and are not production throughput estimates. The source establishes the cause: the new condition falls through to return Ok(true) for non-UTC values, activating the full array.keys_iter() scan at lines 699–702. array_with_timezone preserves dictionary encoding, making this reachable through TimestampTruncExpr.
There was a problem hiding this comment.
Fixed in 9ade55e. SECOND and MILLISECOND now run the values-only range check before the non-UTC guard, so in check mode a dictionary keeps key masking only when one of its values is near the lower bound, whatever the timezone.
With your shape (32 distinct values, Asia/Tokyo, check mode, the kernel at opt-level 3 over debug dependencies), the call drops from 18.4 µs to 2.4 µs at 8,192 keys and from 132 µs to 2.4 µs at 65,536 keys, the same as UTC and the wrapping mode. test_timestamp_trunc_dictionary_fine_mask_decision covers the decision for no timezone, UTC and Asia/Tokyo in both modes.
… dictionaries SECOND and MILLISECOND ignore the timezone, so when they do not wrap they need key masking only for a value near the lower bound. Previously a non-UTC dictionary fell through to the unconditional non-UTC masking branch and scanned every key.
| // smallest timestamp. Spark before 4.2.0 wraps the Long subtraction; 4.2.0 (SPARK-56663) | ||
| // raises `long overflow` like every other unit. This tracks the Spark version Comet is | ||
| // running against. | ||
| bool wrap_second_millisecond_overflow = 4; |
There was a problem hiding this comment.
This is to track the Spark version Comet is running against?
There was a problem hiding this comment.
If so, why it's not named as spark_420_plus?
There was a problem hiding this comment.
I think it is better to name the behavior rather than a Spark version. Some companies maintain forks of Spark and backport fixes from newer OSS Spark versions. These companies can leverate the shim layer to say which of their versions have specific behavior. I have also seen in the past (although rarely) that there are functional differences between different minor releases.
| /// reallocating, and parsed once into a `chrono::TimeZone` per batch. | ||
| timezone: Arc<str>, | ||
| /// Whether SECOND and MILLISECOND truncation wraps below the smallest timestamp, as Spark | ||
| /// did before 4.2.0, instead of raising `long overflow` as 4.2.0 does (SPARK-56663). |
There was a problem hiding this comment.
Is long overflow the behavior for 4.2.0+? Should we cover future Spark version in such comments?
There was a problem hiding this comment.
I don't think we can predict the behavior of future Spark versions
Which issue does this PR close?
Closes #6739.
Closes #6744.
Rationale for this change
Spark 4.2.0 raises
long overflowwhen SECOND or MILLISECOND truncation falls below the smallest timestamp, like every other unit (SPARK-56663). Spark 4.1 and earlier wrap the Long subtraction instead. The fixture added in #5956 expects the wrap, sotrunc_timestamp.sqlfails on the Spark 4.2 profile. Comet's native kernel also wraps on every version, so on Spark 4.2 Comet returns a timestamp in the year 294247 where Spark raises.The pull request tier only runs Spark 4.1, so the failure first showed up on #6664, which runs every profile. It also turns the 4.2 nightly red.
What changes are included in this PR?
TruncTimestampgets awrap_second_millisecond_overflowproto field. The serde sets it to!isSpark42Plus, the same wayMode.normalize_neg_zerotracks the running Spark version. The proto default,false, is the 4.2 behavior.Long.MinValuestill never raises and other dictionaries are still truncated through their values alone.trunc_timestamp.sqlinto a pair of version-gated fixtures.trunc_timestamp_fine_overflow.sqlcovers Spark 4.1 and earlier, where both Spark and Comet wrap.trunc_timestamp_fine_overflow_spark42.sqlcovers 4.2, where both raise. It also runs a nativedate_truncquery over valid input, so the errors cannot come from a fallback to Spark.date_truncexpression audit.How are these changes tested?
test_timestamp_trunc_fine_spark42_overflowcovers the error, the first representable boundary and NULLs for both units across timezones. The dictionary unused-extremes test and the poisoned-NULL test now run in both modes.test_timestamp_trunc_dictionary_fine_mask_decisionchecks that a dictionary of safe values skips key masking in either mode and in any timezone, and that one holding a value near the lower bound is masked only when the units fail.CometSqlFileTestSuite trunc_timestampandCometTemporalExpressionSuite date_truncpass locally on Spark 4.2, 4.1 and 3.5, 27 tests each.trunc_timestamp_fine_overflow_spark42.sqlfails on Spark 4.2 with "Expected Comet to throw an error matching 'long overflow' but query succeeded".