Skip to content
Merged
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
78 changes: 53 additions & 25 deletions docs/source/contributor-guide/timezones.md
Original file line number Diff line number Diff line change
Expand Up @@ -96,30 +96,31 @@ Inside a native plan, every `TimestampType` value must carry exactly `"UTC"`, an
`Timestamp(µs, "UTC")` fails with `Invalid comparison operation`, even though the values are
comparable. DataFusion's `BinaryExpr` does not coerce, and Comet builds binary expressions
directly from Spark's already-typed plan.
- `CASE` and `COALESCE` reconcile differing branch types with casts that are built without a
timezone. Those casts panic when a timestamp branch has to change its label.
- Native datetime kernels decide between wall-clock and instant semantics from the label. A
`TimestampType` value labelled `None` is treated as `TimestampNTZType` and silently loses the
session timezone.

A mislabelled column is easy to miss in tests. `ScanExec` casts every column it imports from the
JVM to its declared type, and the shuffle writer's `SchemaAlignExec` does the same before it
partitions, so a wrong label disappears at the next stage boundary. A test that only projects the
result passes. The label only matters when the result is compared, goes through a `CASE`, or
feeds another native expression.
partitions, so a wrong label disappears at the next stage boundary. `CASE`, `COALESCE` and `IF`
cast a branch whose Arrow type differs from the others to a common type (`coerce_branch` in
`native/spark-expr/src/conditional_funcs/case_when.rs`). For a timestamp that cast only changes the
label, so a mislabelled branch does not fail there either. A test that only projects the result
passes. The label only matters when the result is compared or feeds another native expression.

## How the session timezone reaches native code

Each timezone-aware serde reads the timezone that Spark stamped on the expression
(`expr.timeZoneId`) and serializes it into the expression's protobuf message. That includes `Cast`,
`Hour`, `Minute`, `Second`, `UnixTimestamp`, `TruncTimestamp`, `ToJson`, `ToCsv` and
`ToPrettyString`. Serdes should use this value rather than `SQLConf.get.sessionLocalTimeZone` or
the JVM default, because it is what Spark itself evaluates the expression with.
(`expr.timeZoneId`), converts it with `CometTimeZone.nativeId`, and serializes the result into the
expression's protobuf message. That includes `Cast`, `Hour`, `Minute`, `Second`, `UnixTimestamp`,
`TruncTimestamp`, `ToJson`, `ToCsv` and `ToPrettyString`. Serdes should use this value rather than
`SQLConf.get.sessionLocalTimeZone` or the JVM default, because it is what Spark itself evaluates
the expression with.

When the expression has no timezone, the serdes pass `"UTC"` (`timeZoneId.getOrElse("UTC")`).
Spark does not resolve a timezone-aware expression without one, so in practice the fallback only
applies to casts that do not use the timezone. It cannot simply be removed, though, because the
native side asserts a non-empty timezone even for those casts.
When the expression has no timezone, `nativeId` returns `"UTC"`. Spark does not resolve a
timezone-aware expression without one, so this only happens for casts that do not use the
timezone. Native code reports an empty timezone as an error (`require_timezone` in
`native/spark-expr/src/utils.rs`) rather than guessing one.

On the native side, `array_with_timezone` in `native/spark-expr/src/utils.rs` is the common entry
point:
Expand All @@ -145,11 +146,19 @@ runs natively only in UTC sessions.

### Parsing timezone IDs

Native code parses timezone IDs with arrow's `Tz::from_str`. It accepts IANA names such as
`America/Los_Angeles`, `Etc/UTC` and `UTC`, and fixed offsets written as `+HH`, `+HHMM` or
`+HH:MM`. Spark resolves IDs with `ZoneId.of(id, ZoneId.SHORT_IDS)`, which also accepts `Z`,
offsets such as `+8` and `+08:00:00`, prefixed offsets such as `GMT+8`, and short IDs such as
`PST`. Code that takes a fast path for UTC should compare the ID against a fixed list of UTC
Spark resolves session timezone IDs with `ZoneId.of(id, ZoneId.SHORT_IDS)`, which accepts forms
such as `Z`, offsets like `+8` and `+08:00:00`, prefixed offsets like `GMT+8`, and short IDs like
`PST`. Native code parses timezone IDs with arrow's `Tz::from_str`, which accepts only IANA names
such as `America/Los_Angeles` and fixed offsets written as `+HH`, `+HHMM` or `+HH:MM`.

`CometTimeZone.nativeId` in `spark/src/main/scala/org/apache/comet/serde/CometTimeZone.scala`
bridges the two. It normalizes the zone first, so an ID whose offset never changes becomes `UTC`
for a zero offset, which covers `Z`, `GMT` and `Etc/UTC`, or `+HH:MM` otherwise. A short ID becomes
its region. An offset with seconds, such as `+05:45:30`, has no native spelling, so `nativeId`
returns `None`. The serde then reports the expression as unsupported through
`CometTimeZone.supportLevel`, and it runs in the codegen dispatcher or falls back to Spark.

Code that takes a fast path for UTC should still compare the ID against a fixed list of UTC
aliases, and send everything else down the general path. The list in `extract_date_part.rs` is an
example.

Expand All @@ -159,7 +168,10 @@ Spark converts between instants and local time using the JVM's timezone rules (`
code uses chrono-tz, which compiles its own copy of the IANA database into libcomet
(`chrono_tz::IANA_TZDB_VERSION`), and its precomputed DST transitions end around 2100. The two can
disagree, both for zones whose rules changed between the two database versions and for far-future
timestamps.
timestamps. When the library loads, `NativeBase` compares the native version (`getTzdataVersion`)
with the JVM's and logs a warning if they differ. The user guide describes the effect under
"Timezone Database Versions" on the datetime
[expression compatibility](../user-guide/latest/compatibility/expressions/index.md) page.

## Scans

Expand Down Expand Up @@ -195,20 +207,28 @@ For Iceberg, iceberg-rust labels `timestamptz` columns `Timestamp(Microsecond, "
Iceberg scan adapts its batches to the Spark schema, which relabels them `"UTC"`. Iceberg's
partition transforms (`years`, `months`, `days` and `hours`) are defined in UTC.

The native CSV scan, which is only enabled for testing, parses a timestamp without an offset as
UTC. Spark parses it in the CSV `timeZone` option, which defaults to the session timezone, so
`CometScanRule` falls back when the read schema has a `TimestampType` column and that timezone is
not UTC.

## The codegen dispatcher

Several timezone-dependent expressions have no native implementation that is compatible in every
session, including `date_trunc`, `date_format`, `from_unixtime`, `from_utc_timestamp`,
`to_utc_timestamp`, `make_timestamp` and `to_timestamp`. Their serdes route the cases the native
path cannot handle through the JVM codegen dispatcher. The dispatcher runs Spark's generated code
with the `timeZoneId` stamped on the expression, so the results match Spark, and its Arrow output
is labelled `"UTC"` like everything else.
path cannot handle through the JVM codegen dispatcher, as they do for a session timezone that
`CometTimeZone.nativeId` cannot express. The dispatcher runs Spark's generated code with the
`timeZoneId` stamped on the expression, so the results match Spark, and its Arrow output is
labelled `"UTC"` like everything else.

## Guidelines

- Never read `TimeZone.getDefault`, `ZoneId.systemDefault` or the host's local time in Comet code.
Spark's semantics come from the timezone stamped on the expression. The JVM default only matters
as the default value of `spark.sql.session.timeZone`.
- Pass a timezone to native code only through `CometTimeZone.nativeId`, and return
`CometTimeZone.supportLevel` from `getSupportLevel` when it gives `None`.
- Never label a `TimestampType` value with the session timezone, and never return
`Timestamp(_, None)` for a `TimestampType` result.
- Don't assume `"UTC"` is the only UTC session timezone. `Etc/UTC` is the common default, and a
Expand All @@ -221,13 +241,21 @@ is labelled `"UTC"` like everything else.
- Run in several session timezones: `UTC`, `Etc/UTC`, a zone with DST such as
`America/Los_Angeles`, and a zone with a half-hour offset such as `Asia/Kolkata`. SQL file tests
can use `-- ConfigMatrix: spark.sql.session.timeZone=...` (see
[Comet SQL Tests](sql-file-tests.md)).
- Use the result, rather than only projecting it. Compare it with another timestamp, put it in a
`CASE`, and feed it to `hour` or a cast to string. A wrong label only shows up there.
[Comet SQL Tests](sql-file-tests.md)). `session_timezone_ids.sql` runs the forms that
`CometTimeZone.nativeId` rewrites, such as `GMT+8`, `Z` and `PST`.
- Use the result, rather than only projecting it. Compare it with another timestamp, and feed it to
`hour` or a cast to string. A wrong label only shows up there.
- Include timestamps around DST transitions, before the epoch, and before 1900, when many zones
used local mean time offsets.
- Avoid dates in zones whose rules changed in recent tzdata releases unless the test needs them. A
JDK with older timezone data gives Spark different answers there, and the warning from
`NativeBase` shows when the two versions differ.
- Build inputs from Parquet tables rather than `VALUES` lists. The optimizer evaluates a
projection over `VALUES` itself, so Comet never runs the expression.
- When collecting timestamps in Scala tests, set `spark.sql.datetime.java8API.enabled=true`, or
cast to strings. `java.sql.Timestamp` conversion goes through the JVM's default timezone and the
hybrid calendar.
- Check that the expression ran natively, not only that the plan has Comet operators. The codegen
dispatcher returns Spark's answer from inside a Comet operator, and a dispatched expression takes
its whole subtree with it, including children that have a native implementation. `expect_native`
in SQL file tests and `checkSparkAnswerAndImpl` in Scala tests check which one ran.