Repository navigation
feat: improve Spark from_utc_timestamp compatibility - #25979
nat-openai wants to merge 1 commit into
Conversation
|
@sunchao, if this seems reasonable, can you enable the CI workflows? |
The timezone spellings accepted by Spark are different than those accepted by Arrow's parser. As a result, `from_utc_timestamp` in comet isn't safe to use, since it would fail on things like Java short IDs, prefixed offsets with second precision, and SystemV-style IDs. Additionally, Arrow will accept offsets beyond the Spark limit of 18 hours. This change improves compatibility by parsing timezone IDs with the rules implemented by Spark. As part of doing this, I added SparkFromUtcTimestampExpr so that computed timezones evaluated for rows with non-null timestamps, avoiding failing when a timezone is invalid but wouldn't be required anyway. (Full compatibility with spark behaviour seems quite complex and I haven't attempted it here.) It also adds more documentation recording remaining differences: - timezone validation and evaluation order can differ from Spark - chrono accepts a narrower calendar range than Spark implemented range - chrono-tz doesn't currently perform DST transitions past 2099 This does not achieve full compatibility with Spark due to those issues, but it makes the region of incompatibility a lot smaller.
3244992 to
7594f11
Compare
ajsquared
left a comment
There was a problem hiding this comment.
Reviewed the current changes and feedback; no independently confirmed unresolved P1 remains.
|
@nat-openai just enabled the CI run |
sunchao
left a comment
There was a problem hiding this comment.
LGTM as a partial Spark compatibility improvement. Reviewed 7594f111a214f619c148b98df5b0bc367f75c845 against merge base 3310fd59a743dad59329779451278518ac32d90a; no blocking correctness defect found. The inline suggestions are nonblocking. This does not establish full Spark compatibility.
Summary
-
Prior state and problem: The existing UDF delegated timezone parsing to Arrow, rejecting valid Spark spellings such as
PST,GMT+1, second-resolution offsets, and SystemV identifiers, while accepting offsets beyond Spark's ±18-hour limit. Base/head component checks confirmed thatPSTandGMT+1change from errors to successful conversions. -
Design approach: The private
SparkTimeZoneenum separates fixed offsets, regional timezones, and Java's legacy SystemV rules. Regional conversion continues using Arrow/Chrono. The separateSparkFromUtcTimestampExprcontrols conditional evaluation of timezone expressions, while both entry points share the conversion kernel. -
Correctness: Offset normalization, signs, seconds, ±18-hour boundaries, Java short-ID mappings, and SystemV transitions checked out. Focused checks passed for nulls, scalar/array combinations, sliced validity buffers, empty batches, timestamp precision, and timezone metadata. The shared helper's checked addition turns calendar-overflow panics into errors. No introduced P1/P2 found.
-
Compatibility analysis: Important limits remain. The new evaluator deliberately differs from Spark's generated execution: an invalid literal timezone can be skipped when every timestamp is null, and a null timezone does not suppress timestamp evaluation. Chrono's calendar range is narrower. Regional DST rules also remain incorrect for some future dates: in July 2100, this implementation applies UTC−5 to New York while Spark applies UTC−4. The base UDF reproduces the same future-DST limitation. This PR alone therefore does not justify removing Comet's incompatibility guard.
-
Key design decisions: Preserving offset seconds is necessary. Mapping
EST,HST, andMSTto fixed offsets while mapping other short IDs to regions matches Java. Keeping SystemV rules separate avoids substituting modern regional DST rules. The one-entry timezone cache avoids reparsing constants/consecutive repeated values without persistent state. -
Implementation sketch: Evaluate timestamps; skip timezone evaluation for empty/all-null inputs; read literal/column timezones directly or evaluate computed timezones only on selected rows; parse/cache each needed timezone; shift the timestamp while preserving its type and metadata. The UDF shares the conversion steps but receives already-evaluated arguments.
-
Performance: The direct literal/column fast path and parsing cache are sensible. For computed timezone expressions with mixed-null timestamps, the default
evaluate_selectionfilters every input column, including unrelated columns; an isolated helper probe confirmed scaling with batch width. The benchmark currently exercises only direct timezone columns. Scalar timezone strings also still expand into full arrays throughmake_scalar_function: the cache removes repeated parsing but not that allocation. Scalar expansion is pre-existing overhead, not a new regression. Add representative benchmarks before introducing broader caching or specialized SystemV optimizations. -
Design: The approach is solid and reuses the regional database and selection machinery. The important integration boundary is that ordinary SQL still uses the eager UDF; consumers must explicitly construct the new physical expression to obtain conditional evaluation. A runtime-table probe reproduced the difference using a computed timezone with an invalid cast on null timestamp rows. A planner/consumer example would clarify the API's intended use.
-
Abstraction & complexity: The three timezone variants represent distinct policies and are appropriately scoped. A custom physical expression is justified because an ordinary UDF cannot control evaluation of arguments it has already received. The small last-value cache is proportionate; a generic cache framework would need workload evidence. The bespoke SystemV rules create maintenance responsibility, but their limited scope and transition coverage make that manageable.
-
Behavioral changes worth calling out: More Spark timezone spellings succeed; offsets outside ±18 hours and
ROCare rejected; calendar-boundary adjustment returns an error instead of panicking. Conditional evaluation is opt-in through the physical expression.to_utc_timestampretains its existing parser, so the reverse function does not gain equivalent spelling support here. -
Suggested improvements: Add computed/scalar timezone benchmark cases, including wide batches and sparse nulls. Add a consumer/planner example and integration tests when wiring in the physical expression. As a minor cleanup, have
fmt_sqlcall the children's SQL formatters rather than theirDisplayimplementations, which can emit column indices such astimestamp@0.
Validation
cargo test --locked -p datafusion-spark --lib function::datetime -- --nocapture: 52 passed, including all 15 tests in the changedfrom_utc_timestampand new timezone modules.cargo test --locked -p datafusion-functions --lib datetime::to_local_time -- --nocapture: 4 passed.- Independent SQL/physical component probes using unchanged PR source and its compiled dependencies passed. The base-UDF control reused the head's shared helper and dependencies; it was not a full base build.
- 30 actual Spark 4.1.3 SQL probes checked generated and interpreted execution with UTC session timezone and ANSI enabled. Source comparisons also covered Spark 3.5.8, 4.0.4, and 4.2.0; those versions were not runtime-validated.
- Exact extracted Rust offset-parser functions matched Java with Spark normalization on 258,744 cases. Extracted calendar/helpers passed 6,611,852 checks against JDK 21/tzdb 2026b. These are helper-level differential checks, not full-engine tests.
- Full local SLT/workspace tests and release benchmarks were not run. The selection-cost probe used debug dependencies on a shared machine and is not a release or base/head performance benchmark.
- CI snapshot immediately before submission: 36 successful, 5 running, and 3 skipped checks. Approval reflects the reviewed code and focused validation; full CI completion remains outstanding.
| let rows = 8192; | ||
| let expr = SparkFromUtcTimestampExpr::new( | ||
| Arc::new(Column::new("timestamp", 0)), | ||
| Arc::new(Column::new("timezone", 1)), |
There was a problem hiding this comment.
Nonblocking: could we also benchmark a computed timezone expression, with sparse and mixed null timestamps, across these batch widths? This Column takes the fast path and bypasses evaluate_selection, so it does not exercise the branch that filters every input column. A computed-expression case would measure the remaining allocation/width tradeoff; a literal timezone case would also cover scalar broadcasting.
| } | ||
|
|
||
| impl SparkFromUtcTimestampExpr { | ||
| /// Construct an expression that skips computed timezone evaluation for null timestamps. |
There was a problem hiding this comment.
Nonblocking documentation suggestion: a short planner/consumer example would help show where to construct this expression instead of the registered UDF, and that its children need to be coerced before construction. Ordinary SQL still uses the eager UDF, while this constructor simply stores the supplied physical children. An integration example would make the boundary of the conditional-evaluation support clearer.
|
|
||
| impl PhysicalExpr for SparkFromUtcTimestampExpr { | ||
| fn fmt_sql(&self, f: &mut Formatter<'_>) -> std::fmt::Result { | ||
| Display::fmt(self, f) |
There was a problem hiding this comment.
Nit: could fmt_sql format each child through its fmt_sql implementation rather than delegate to Display? For example, column Display produces timestamp@0, whereas its SQL formatter produces timestamp. This would keep SQL-formatted output consistent with other physical expressions.
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25979 +/- ##
==========================================
+ Coverage 82.77% 82.79% +0.01%
==========================================
Files 1147 1148 +1
Lines 450926 451621 +695
Branches 450926 451621 +695
==========================================
+ Hits 373275 373915 +640
- Misses 54927 54956 +29
- Partials 22724 22750 +26 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
Which issue does this PR close?
This addresses the major incompatibility described in apache/datafusion-comet#4654 for
CometFromUTCTimestamp. It does not close that issue.It's a follow-up to #19879.
Rationale for this change
The timezone spellings accepted by Spark are different than those accepted by Arrow's parser. As a result,
from_utc_timestampin comet isn't safe to use, since it would fail on things like Java short IDs, prefixed offsets with second precision, and SystemV-style IDs. Additionally, Arrow will accept offsets beyond the Spark limit of 18 hours.What changes are included in this PR?
This change improves compatibility by parsing timezone IDs with the rules implemented by Spark. As part of doing this, I added SparkFromUtcTimestampExpr so that computed timezones evaluated for rows with non-null timestamps, avoiding failing when a timezone is invalid but wouldn't be required anyway. (Full compatibility with spark behaviour seems quite complex and I haven't attempted it here.)
It also adds more documentation recording remaining differences:
This does not achieve full compatibility with Spark due to those issues, but it makes the region of incompatibility a lot smaller.
If this approach isn't the one that folks would prefer, please let me know if there are different avenues you'd rather I explore!
Finally, it's worth mentioning that this change was LLM-assisted.
What is the testing strategy for this PR?
Added unit tests in
datafusion/spark/src/function/datetime/timezone.rscovering accepted spelling, the 18 hour boundary, and SystemV custom behaviour around DST transitions and exception.Added SQL logic tests in
datafusion/sqllogictest/test_files/spark/datetime/from_utc_timestamp.sltcovering the accepted spellings.Are there any user-facing changes?
Yes. The existing Spark
from_utc_timestampUDF now accepts additional Spark timezone spellings. It also rejects offsets outside Spark's 18-hour limit.This won't affect most users since they'd need to have opted into the existing incompatible implementation.