Timezone Handling#
This page describes how Comet represents timestamps and applies the Spark session timezone. It is aimed at contributors working on datetime expressions, casts, scans, or any native code that produces or consumes timestamp columns. For user-facing differences from Spark, see the datetime and cast sections of the Expression Compatibility pages, and the scan compatibility notes. Known bugs are tracked in #6335.
The short version: Comet never converts timestamp values at the JVM/native boundary. The session timezone is not part of a value. It travels with each timezone-aware expression and is applied inside the kernel that needs it.
How Spark models time#
Spark type |
Stored value |
Timezone |
|---|---|---|
|
Microseconds since the Unix epoch: an instant |
The session timezone, when converting to or from local time |
|
Microseconds since the epoch of a wall-clock date and time |
None |
|
Days since the epoch |
None |
The session timezone is spark.sql.session.timeZone. Its default is the JVM’s default timezone
when the session is created. On Linux hosts and containers whose /etc/localtime points at
Etc/UTC, which is the default on Ubuntu and Debian images, that is Etc/UTC, not UTC.
During analysis, Spark’s ResolveTimeZone rule stamps the session timezone onto every
TimeZoneAwareExpression, and such an expression is not resolved until it has one. Cast is the
exception. It needs a timezone only when Cast.needsTimeZone(from, to) is true: string to or from
timestamp, timestamp to or from date, timestamp to or from TIMESTAMP_NTZ, and those same pairs
nested in arrays, maps and structs. A cast that Spark or Comet creates after analysis can therefore
legitimately arrive with timeZoneId = None.
Spark uses the session timezone to:
convert a
TimestampTypevalue to or from a string, a date, or aTimestampNTZTypevalueextract local fields:
hour,minuteandsecond, and the date fields through an implicit cast to datetruncate (
date_trunc), format (date_format,from_unixtime) and parse (to_timestamp,unix_timestampof a string or a date)evaluate
from_utc_timestamp,to_utc_timestampandconvert_timezoneadd calendar intervals, whose day and month units follow local time
It does not use the timezone for comparisons, sorting, hashing, unix_timestamp of a timestamp,
or TimestampNTZType values, apart from converting them to TimestampType. For a
TimestampNTZType input, the field extractors use UTC regardless of the session timezone
(zoneIdForType in Spark’s datetimeExpressions.scala).
How Comet represents timestamps#
Spark type |
Arrow type in native code |
|---|---|
|
|
|
|
|
|
The "UTC" is a label, not a conversion. A TimestampType value is already a UTC instant, so
Comet passes the raw microseconds across the boundary in both directions without shifting them.
The label is set in these places:
to_arrow_datatypeinnative/core/src/execution/serde.rs, for every serialized type: scan schemas, expression types and literals.CometArrowStream.NATIVE_TIMEZONEinCometNativeArrowSource.scala, which every JVM producer uses when it exports Spark data to native code. That includesCometSparkToColumnarExec,CometLocalTableScanExec, the in-memory cache serializer, and the inputs to shuffles and writes.The output schema of the codegen dispatcher, in
CometBatchKernelCodegenOutput.scala.
On the way back, Utils.fromArrowType maps any microsecond timestamp that has a timezone to
TimestampType, and CometVector reads the values unchanged.
The invariant#
Inside a native plan, every TimestampType value must carry exactly "UTC", and every
TimestampNTZType value no timezone at all. Several things depend on this:
Arrow’s comparison kernels require identical types.
Timestamp(µs, "Etc/UTC")compared withTimestamp(µs, "UTC")fails withInvalid comparison operation, even though the values are comparable. DataFusion’sBinaryExprdoes not coerce, and Comet builds binary expressions directly from Spark’s already-typed plan.CASEandCOALESCEreconcile 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
TimestampTypevalue labelledNoneis treated asTimestampNTZTypeand 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.
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.
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.
On the native side, array_with_timezone in native/spark-expr/src/utils.rs is the common entry
point:
A
TimestampTypeinput is relabelled with the session timezone. For casts to strings and dates, the values are also shifted to local time and the label is dropped.A
TimestampNTZTypeinput is left alone, except when it is cast toTimestampType. In that case the wall-clock value is resolved in the session timezone.
DST transitions follow Java (resolve_local_datetime). An ambiguous local time takes the earlier
offset. A local time inside a gap resolves with the offset that applied before the transition,
which gives the same instant as Java’s atZone.
A native kernel that produces a TimestampType value must do its local-time work in the session
timezone and then return the result labelled "UTC". That applies both to the declared type, in
data_type() or return_type(), and to the arrays it builds.
DataFusion’s own datetime functions take their timezone from the argument’s label, or from
datafusion.execution.time_zone, which Comet leaves unset. Wiring one of them in for a
timezone-aware Spark expression therefore evaluates it in UTC, unless the session timezone is
passed explicitly. For example, to_char formats in the input’s label, which is why date_format
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
aliases, and send everything else down the general path. The list in extract_date_part.rs is an
example.
Timezone rules#
Spark converts between instants and local time using the JVM’s timezone rules (tzdb.dat). Native
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.
Scans#
Parquet stores timestamps either as INT64 annotated TIMESTAMP(MICROS or MILLIS, isAdjustedToUTC),
or as legacy INT96. Comet’s native scan follows Spark’s vectorized reader, which never shifts a
value by a timezone:
Parquet column |
Spark type |
Comet |
|---|---|---|
|
|
Read as is. Milliseconds are scaled to microseconds. |
|
|
Relabelled without shifting, as Spark does |
|
|
Read as is |
|
|
Rejected on Spark 3.x (SPARK-36182). Relabelled without shifting on Spark 4.0+ (SPARK-47447). |
|
|
Coerced to microseconds and labelled |
The INT96 coercion is configured by coerce_int96 and coerce_int96_tz in parquet_exec.rs. An
INT96 column read as TimestampNTZType follows the isAdjustedToUTC=true row.
The timestamp-to-timestamp adaptations go through CometCastColumnExpr and
parquet_convert_array in native/core/src/parquet/, not through Spark’s Cast. A Spark cast
between TimestampType and TimestampNTZType would apply the session timezone, and the Parquet
reader does not.
Two Spark settings are not handled natively:
Comet disables itself for a session with
spark.sql.parquet.int96TimestampConversion=true, which shiftsINT96values written by Impala.Comet does not rebase dates and timestamps that were written with the legacy hybrid calendar. See #5010. Spark’s rebase of legacy timestamps is itself timezone-dependent.
For Iceberg, iceberg-rust labels timestamptz columns Timestamp(Microsecond, "+00:00"). The
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 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.
Guidelines#
Never read
TimeZone.getDefault,ZoneId.systemDefaultor 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 ofspark.sql.session.timeZone.Never label a
TimestampTypevalue with the session timezone, and never returnTimestamp(_, None)for aTimestampTyperesult.Don’t assume
"UTC"is the only UTC session timezone.Etc/UTCis the common default, and a native path gated on “the session is UTC” must still return"UTC"-labelled output there.Don’t apply a timezone to a
TimestampNTZTypevalue, except when converting it toTimestampType.
Testing timezone-sensitive code#
Run in several session timezones:
UTC,Etc/UTC, a zone with DST such asAmerica/Los_Angeles, and a zone with a half-hour offset such asAsia/Kolkata. SQL file tests can use-- ConfigMatrix: spark.sql.session.timeZone=...(see Comet SQL Tests).Use the result, rather than only projecting it. Compare it with another timestamp, put it in a
CASE, and feed it tohouror 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.
Build inputs from Parquet tables rather than
VALUESlists. The optimizer evaluates a projection overVALUESitself, 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.Timestampconversion goes through the JVM’s default timezone and the hybrid calendar.