Skip to content

Datetime rebase: track the documented scan limitation, and spark.comet.exceptionOnDatetimeRebase is dead code #5010

Description

@andygrove

Describe the bug

Updated after filing: the core limitation here is already documented in the user guide, in compatibility/scans.md under "The following limitation may produce incorrect results without falling back to Spark":

No support for datetime rebasing. When reading Parquet files containing dates or timestamps written before Spark 3.0 (which used a hybrid Julian/Gregorian calendar), dates/timestamps will be read as if they were written using the Proleptic Gregorian calendar. This may produce incorrect results for dates before October 15, 1582.

This issue is retained as the tracking issue for that limitation, which the doc does not currently link to. #1638 previously tracked the same gap for native_iceberg_compat and was closed as obsolete when that reader was removed; the gap still exists in the current native scan. The INT96 variant was filed separately as #5011 and closed as a duplicate of this issue.

Two aspects are not covered by the documented text and are the actionable part of this issue.

1. spark.comet.exceptionOnDatetimeRebase is dead code

The config is declared at CometConf.scala:695, is not .internal(), and is in CATEGORY_EXEC, so it publishes to the user-facing configs page. Its doc string promises:

When this is true, Comet will throw exceptions when seeing these dates/timestamps that were written by Spark version before 3.0.

It is not read anywhere else in the tree. SparkParquetOptions::use_legacy_date_timestamp_or_ntz (native/core/src/parquet/parquet_support.rs:78), which appears to be its intended destination, is only ever initialized to false and never consulted. A user who sets the config to true specifically to be protected from this limitation gets silence and wrong results instead.

Either wire the config up so it raises on legacy-calendar data, or remove it so it stops advertising a guarantee that does not exist.

2. The documented scope is narrower than the actual behavior

  • The doc says "written before Spark 3.0". Any current Spark writes legacy-calendar data when spark.sql.parquet.datetimeRebaseModeInWrite=LEGACY is set, and Comet mis-reads those files too. The repro below uses Spark 4.1 as the writer.
  • The doc implies visibly-wrong dates. In practice filters and aggregates over the affected column silently return the wrong answer, which is harder to notice than an obviously-shifted date.
  • spark.sql.parquet.datetimeRebaseModeInRead / spark.sql.legacy.parquet.datetimeRebaseModeInRead are not mentioned. Comet reads neither, and also ignores the org.apache.spark.legacyDateTime file metadata that takes precedence over them in Spark.

Steps to reproduce

val path = "/tmp/legacy_dates"

// Written by Spark 4.1, not by a pre-3.0 Spark.
spark.conf.set("spark.sql.parquet.datetimeRebaseModeInWrite", "LEGACY")
spark.sql("SELECT cast(s as date) as d FROM VALUES ('1000-01-01'),('1990-01-01') AS v(s)")
  .write.mode("overwrite").parquet(path)
spark.conf.unset("spark.sql.parquet.datetimeRebaseModeInWrite")

spark.conf.set("spark.comet.enabled", "false")
spark.read.parquet(path).show()
// 1000-01-01
// 1990-01-01

spark.conf.set("spark.comet.enabled", "true")
spark.read.parquet(path).show()
// 1000-01-06   <-- shifted by the Julian/Gregorian offset
// 1990-01-01

Setting spark.comet.exceptionOnDatetimeRebase=true changes nothing.

Silent divergence in filters and aggregates, not just projection:

SELECT count(*) FROM parquet.`/tmp/legacy_dates` WHERE d = date'1000-01-01'
-- spark: 1
-- comet: 0

SELECT min(d), max(d) FROM parquet.`/tmp/legacy_dates`
-- spark: 1000-01-01, 1990-01-01
-- comet: 1000-01-06, 1990-01-01

TIMESTAMP_MICROS written the same way behaves the same (1000-01-01 12:34:56 reads back as 1000-01-06 12:41:58). spark.comet.scan.enabled=false makes every case above match, isolating this to the native scan.

Expected behavior

For the underlying limitation, either of:

  1. the native scan rebases legacy-calendar dates/timestamps as Spark does, honouring the file metadata first and spark.sql.parquet.datetimeRebaseModeInRead when it is absent; or
  2. Comet falls back to Spark for scans of files carrying org.apache.spark.legacyDateTime and for non-default values of the read config.

Independently of which is chosen, spark.comet.exceptionOnDatetimeRebase should either work or be removed, and scans.md should be widened as described above.

Additional context

Activity

  1. added
    bugSomething isn't working
    priority:criticalData corruption, silent wrong results, security issues
    area:scanParquet scan / data reading
    on Jul 23, 2026
  2. changed the title [-]Native Parquet scan ignores legacy-calendar rebase, returning silently wrong dates/timestamps[/-] [+]Datetime rebase: track the documented scan limitation, and spark.comet.exceptionOnDatetimeRebase is dead code[/+] on Jul 23, 2026
  3. added
    priority:mediumFunctional bugs, performance regressions, broken features
    documentationImprovements or additions to documentation
    and removed
    priority:criticalData corruption, silent wrong results, security issues
    on Jul 23, 2026
  4. peterxcli commented on Jul 24, 2026

    @peterxcli
    Member

    take

  5. andygrove commented on Aug 1, 2026

    @andygrove
    MemberAuthor

    Re-verified on apache/main @ ba21f02de (default Maven profile, Spark 4.1.2, macOS aarch64, debug
    build) while re-running the #4180 audit. Still reproduces, unchanged.

    Two things the original report did not pin down, both now measured:

    All three read-mode values behave identically. I tested
    spark.sql.parquet.datetimeRebaseModeInRead at LEGACY, CORRECTED and EXCEPTION, and also at
    the session default, against a file written by Spark 4.1 with
    datetimeRebaseModeInWrite=LEGACY. Comet returns 1000-01-06 in every one of the four, where
    Spark returns 1000-01-01. The config is not merely mishandled for some values, it is never
    consulted. spark.sql.parquet.int96RebaseModeInRead behaves the same across the same four cases
    (1000-01-01 12:34:56 reads back as 1000-01-06 12:41:58).

    The predicate case is silent, not just visibly shifted.

    SELECT count(*) FROM legacy_dates WHERE d = date'1000-01-01'
    -- spark: 1
    -- comet: 0

    Carve-out: part 1 of this issue, spark.comet.exceptionOnDatetimeRebase being dead code, is now
    tracked separately as #5195. It has a self-contained fix (wire the config up, or remove it) that
    does not depend on rebase support landing, and it was blocking this issue from reading as what it
    is. Part 2, the documentation scope, has since been addressed: compatibility/scans.md now
    describes the wider behaviour, names datetimeRebaseModeInRead and the
    org.apache.spark.legacyDateTime metadata, and links here.

    That leaves this issue as the underlying correctness gap only: either rebase legacy-calendar
    dates and timestamps in the native scan as Spark does, or fall back the affected scan to Spark.

  6. andygrove commented on Aug 2, 2026

    @andygrove
    MemberAuthor

    Findings from attempting a fix (#5202, closed)

    I closed #5202 ("fail rather than silently return unrebased datetimes from Parquet") because the
    problem is substantially more complex than the issue text implies. #5221 now only removes the dead
    spark.comet.exceptionOnDatetimeRebase config for 1.0.0, so the correctness gap tracked here is
    untouched and still open.

    Recording what the attempt turned up, so the next person picking this up (or reviewing #5048)
    starts from it rather than rediscovering it. All of the below is verified against Spark's own
    sources, not inferred.


    A. Detection is much narrower than "does the file carry a legacy marker"

    A1. The footer rule is "provably Proleptic Gregorian", not "provably legacy". Spark's
    DataSourceUtils.getRebaseSpec
    resolves the policy as: if org.apache.spark.version is present, LEGACY when
    version < minVersion || <marker key> != null, else CORRECTED; if the version key is absent
    entirely, fall through to the read-mode config.

    The absent-version case is not hypothetical and is easy to get wrong: Spark 2.4.5 and earlier
    wrote no org.apache.spark.version key at all
    , so the canonical legacy files carry no metadata
    hint whatsoever. Any detector keyed on "marker present" or "version < 3.0.0" clears exactly the
    files that need rebasing most. Credit to @peterxcli — #5048 got this right and it is what caught
    the same bug in my first commit.

    Note Spark compares with version < minVersion as a string, not a semantic version compare. A
    native or JVM reimplementation should replicate the string comparison rather than "fixing" it, or
    it will diverge on odd version strings.

    A2. Two markers, two thresholds, tracked independently.
    org.apache.spark.legacyDateTime with minVersion 3.0.0 for date/TIMESTAMP_MICROS/
    TIMESTAMP_MILLIS, and org.apache.spark.legacyINT96 with minVersion 3.1.0 for INT96.
    Collapsing them into one flag or one version threshold is wrong in both directions.

    A3. Read modes come from ParquetOptions, not SQLConf. Spark resolves
    datetimeRebaseModeInRead / int96RebaseModeInRead through
    new ParquetOptions(options, conf), so a per-read .option("datetimeRebaseMode", "CORRECTED")
    overrides the session conf. Reading SQLConf directly silently ignores it. (Again @peterxcli's
    catch, from reviewing #5202.)

    A4. The file-level marker says nothing about whether any value is actually affected. Spark
    stamps legacyDateTime on a whole file whenever the write mode was LEGACY, regardless of the
    values — and dates from 1582-10-15 onward rebase to themselves
    (RebaseDateTime.lastSwitchJulianDay == -141427, and julianGregDiffs.last == 0). So a large
    fraction of marked files return perfectly correct results under Comet today. Marker-only detection
    over-triggers heavily: whichever remedy is chosen (raise, fall back, or rebase), it fires on reads
    that need nothing.

    Row-group min/max statistics are the only cheap discriminator, and they are only available where
    the footer is already in hand.

    A5. Statistics cannot save INT96. The Parquet spec gives INT96's 12 bytes no meaningful
    ordering, so writers emit no usable min/max. INT96 columns therefore have to be handled
    conservatively regardless of statistics. Same for any row group with no statistics. An all-null row
    group is clear (no value to rebase), and so is an empty one.

    A6. Thresholds must come from Spark, not from native constants.
    RebaseDateTime.lastSwitchJulianTs is not a constant anyone should re-derive: it is
    rebaseMap.values.map(_.switches.last).max over Spark's per-timezone rebase tables, with a
    require that all diffs after it are zero. Read lastSwitchJulianDay and lastSwitchJulianTs on
    the JVM and send them through the proto, so the thresholds stay exact for whichever Spark version
    is in use. Scaling them to millis/nanos has to round toward the unsafe direction and saturate
    rather than overflow.

    A7. Timestamp rebasing is timezone-dependent, and the timezone is per file. When the resolved
    policy is LEGACY, Spark builds RebaseSpec(LEGACY, Option(lookupFileMeta(SPARK_TIMEZONE_METADATA_KEY)))
    — i.e. it rebases using the writer's org.apache.spark.timeZone, falling back to the session
    timezone. Any actual native rebase implementation (option 1 below) needs the full per-timezone
    tables and per-file timezone resolution, not a single global offset. This is the bulk of the work
    in option 1 and is why it is not a small change.

    A8. The affected-column test must walk nested types. Date/timestamp columns hide inside
    structs, lists, maps, dictionaries and unions, so the walk has to recurse — and it must consider
    only the requested top-level columns (case-insensitively), or reading a modern date column out of
    a file that also holds an ancient one fails for no reason.


    B. Where the check runs decides which remedies are even reachable

    B1. Plan-time detection costs real planning latency. #5048's requiresDatetimeRebase opens
    every selected file's footer on the driver, per query, for any query touching a date or timestamp
    column. Three things compound: it uses selectedPartitions filtered with
    partitionFilters.filterNot(isDynamicPruningFilter), so it is the pre-DPP file set;
    ParquetFileReader.open reads the full footer including all row-group metadata rather than
    ParquetFooterReader.readFooter(..., SKIP_ROW_GROUPS); and it is serial. On a table with thousands
    of files that is significant added planning time and driver heap, paid even when every file is
    fine.

    B2. Read-time detection is free but cannot fall back. In the native reader the footer has
    already been fetched and cached for the read itself, so inspection costs no extra I/O, and it is the
    only point where row-group statistics (A4) are available. But by then the plan is fixed: the only
    available outcome is to raise. This is the core tension — correct-results-via-fallback requires
    paying B1; free detection can only fail the query.

    B3. If done natively, the enforcement point must be the reader factory's get_metadata. Not the
    expression/schema adapter: DataFusion only constructs that when the logical and physical schemas
    differ or a predicate is pushed down, so a plain SELECT d FROM t skips it entirely.


    C. If the remedy is to raise, the error plumbing is not trivial

    C1. It has to surface as Spark's own exception. SparkUpgradeException with
    INCONSISTENT_BEHAVIOR_CROSS_VERSION.READ_ANCIENT_DATETIME. This matters mechanically, not just
    cosmetically: FileScanRDD deliberately rethrows SparkUpgradeException (two sites) instead of
    wrapping it in FAILED_READ_FILE/Encountered error while reading file, because the file is not
    corrupt. Anything that loses the type gets misreported as a corrupt-file error.

    C2. Getting it out of the native reader needs a downcast, not message matching.
    AsyncFileReader::get_metadata can only return a ParquetError, so the typed error has to travel
    boxed in ParquetError::External and be recovered by downcast in the error classifier.

    C3. Unresolved: the message conflicts with Comet's actual remedy. Spark's templated
    READ_ANCIENT_DATETIME message advises setting datetimeRebaseModeInRead to LEGACY/CORRECTED,
    which does not help under Comet — the remedy is to disable Comet for the query. Matching Spark's
    exception type means Comet's real guidance can only travel as the cause; a Comet-specific exception
    makes the guidance primary but is no longer a SparkUpgradeException (see C1). No good answer found.


    D. A raise-only remedy leaves a real divergence

    Under read mode LEGACY on a version-less file, Spark rebases the values and returns them
    correctly. Comet has no rebasing to do it with, so it can only refuse a read Spark completes
    successfully. Only option 1 closes this.


    E. Test assets that already exist, and one test that asserts the bug

    The 15 before_1582_* fixtures are already checked in under
    spark/src/test/resources/test-data/ — date / TIMESTAMP_MICROS / TIMESTAMP_MILLIS / INT96 plain
    / INT96 dict, each written by Spark 2.4.5, 2.4.6 and 3.2.0. That covers every physical encoding and
    both footer paths: the 2.4.x files are the version-less case (A1), the 3.2.0 files the
    marker-stamped case. Any fix should be validated against all 15.

    ParquetReadSuite.scala:1714 ("reading ancient dates before 1582") currently asserts the wrong
    behaviour
    — its comment states "no rebase, no exception". It has to change with any fix here.


    F. Also still outstanding after #5221

    #5221 removes the config from CometConf.scala and the scans.md sentence, but not the dead native
    field SparkParquetOptions::use_legacy_date_timestamp_or_ntz
    (native/core/src/parquet/parquet_support.rs:78, only ever initialized to false at lines 111 and
    127, never consulted). Both #5048 and #5202 removed it; neither has merged, so it survives. Worth
    removing independently.


    Where that leaves the options

    With the costs now measured, restating the choices in the issue body:

    1. Rebase natively. The only option with no divergence (D) and no spurious failures (A4), but
      it needs Spark's per-timezone rebase tables and per-file timezone resolution ported natively
      (A7). Largest change by a wide margin. feat: add native Parquet datetime rebasing #5047 (draft) is the existing attempt.
    2. Fall back to Spark. Always correct, but requires plan-time footer reads (B1) and cannot use
      row-group statistics, so it over-triggers (A4). The plan-time cost is the objection to fix: fall back for Parquet datetime rebasing #5048;
      it could be reduced (SKIP_ROW_GROUPS, post-DPP file set, parallelized) but not eliminated.
    3. Raise. Free to detect and precise thanks to statistics (B2, A4), but fails queries Spark
      answers (D) and needs the error plumbing in section C. This was fix: fail rather than silently return unrebased datetimes from Parquet #5202.

    My own read is that option 3 is not worth shipping on its own — turning silently-wrong results into
    a hard failure is an improvement, but a narrower one than it looks once D is accounted for, and it
    spends the error-plumbing complexity without moving toward option 1. Option 2 with the plan-time
    cost reduced is the pragmatic interim, and option 1 is the real fix. Other opinions welcome —
    cc @peterxcli.

  7. peterxcli commented on Aug 2, 2026

    @peterxcli
    Member

    Option 2 with the plan-time cost reduced is the pragmatic interim

    @andygrove thanks for the deep analysis, I agree and will go with option 2, too. and I think the ROI of option 1 is too low, seems like copy the datetime rebase natively require massive changes.
    Let me revise #5048 first.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:scanParquet scan / data readingbugSomething isn't workingcorrectnessdocumentationImprovements or additions to documentationpriority:mediumFunctional bugs, performance regressions, broken features

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions