From 10d98864f14cfd20393eea175fa1f19ea0ca4769 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Sun, 26 Jul 2026 14:41:35 +0800 Subject: [PATCH 1/9] fix: fall back for Parquet datetime rebasing --- .../user-guide/latest/compatibility/scans.md | 12 +++--- native/core/src/parquet/parquet_support.rs | 4 -- .../scala/org/apache/comet/CometConf.scala | 12 ------ .../comet/parquet/CometParquetUtils.scala | 29 +++++++++++++++ .../apache/comet/rules/CometScanRule.scala | 37 ++++++++++++++++++- .../comet/parquet/ParquetReadSuite.scala | 23 ++++++++++++ 6 files changed, 92 insertions(+), 25 deletions(-) diff --git a/docs/source/user-guide/latest/compatibility/scans.md b/docs/source/user-guide/latest/compatibility/scans.md index aa2be5b3ce8..d67eac0e247 100644 --- a/docs/source/user-guide/latest/compatibility/scans.md +++ b/docs/source/user-guide/latest/compatibility/scans.md @@ -43,19 +43,17 @@ The following features are not supported and cause Comet to fall back to Spark: - No support for `input_file_name()`, `input_file_block_start()`, or `input_file_block_length()` SQL functions. Comet's Parquet scan does not use Spark's `FileScanRDD`, so these functions cannot populate their values. - No support for `ignoreMissingFiles` or `ignoreCorruptFiles` being set to `true` +- Files that require datetime rebasing. Comet falls back to Spark when Parquet metadata or + `spark.sql.parquet.datetimeRebaseModeInRead` / + `spark.sql.parquet.int96RebaseModeInRead` requires legacy-calendar handling. This includes + files written by any Spark version with a corresponding legacy rebase mode, not only files + written before Spark 3.0. See [#5010](https://github.com/apache/datafusion-comet/issues/5010). - `spark.sql.parquet.enableVectorizedReader=false`. Disabling the vectorized reader opts into Spark's parquet-mr semantics (silent overflow, null-on-narrowing), which Comet's native reader does not replicate. By default Comet falls back to Spark in this case. Set `spark.comet.scan.allowDisabledParquetVectorizedReader=true` to opt in to running the Comet Parquet scan regardless. -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. - The following limitations raise an error at scan time rather than falling back to Spark: - Invalid UTF-8 bytes in `STRING` columns. Spark permits arbitrary byte sequences in a `STRING` diff --git a/native/core/src/parquet/parquet_support.rs b/native/core/src/parquet/parquet_support.rs index 2ee1230ed87..4191d876573 100644 --- a/native/core/src/parquet/parquet_support.rs +++ b/native/core/src/parquet/parquet_support.rs @@ -74,8 +74,6 @@ pub struct SparkParquetOptions { pub allow_incompat: bool, /// Support casting unsigned ints to signed ints (used by Parquet SchemaAdapter) pub allow_cast_unsigned_ints: bool, - /// Whether to read dates/timestamps that were written in the legacy hybrid Julian + Gregorian calendar as it is. If false, throw exceptions instead. If the spark type is TimestampNTZ, this should be true. - pub use_legacy_date_timestamp_or_ntz: bool, // Whether schema field names are case sensitive pub case_sensitive: bool, /// SPARK-53535 (Spark 4.1+): when reading a struct whose requested fields are all @@ -108,7 +106,6 @@ impl SparkParquetOptions { timezone: timezone.to_string(), allow_incompat, allow_cast_unsigned_ints: false, - use_legacy_date_timestamp_or_ntz: false, case_sensitive: false, return_null_struct_if_all_fields_missing: true, use_field_id: false, @@ -124,7 +121,6 @@ impl SparkParquetOptions { timezone: "".to_string(), allow_incompat, allow_cast_unsigned_ints: false, - use_legacy_date_timestamp_or_ntz: false, case_sensitive: false, return_null_struct_if_all_fields_missing: true, use_field_id: false, diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index 958a22cc2bf..0fea7ce8759 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -692,18 +692,6 @@ object CometConf extends ShimCometConf { .booleanConf .createWithDefault(false) - val COMET_EXCEPTION_ON_LEGACY_DATE_TIMESTAMP: ConfigEntry[Boolean] = - conf("spark.comet.exceptionOnDatetimeRebase") - .category(CATEGORY_EXEC) - .doc("Whether to throw exception when seeing dates/timestamps from the legacy hybrid " + - "(Julian + Gregorian) calendar. Since Spark 3, dates/timestamps were written according " + - "to the Proleptic Gregorian calendar. When this is true, Comet will " + - "throw exceptions when seeing these dates/timestamps that were written by Spark version " + - "before 3.0. If this is false, these dates/timestamps will be read as if they were " + - "written to the Proleptic Gregorian calendar and will not be rebased.") - .booleanConf - .createWithDefault(false) - val COMET_ENABLE_PARTIAL_HASH_AGGREGATE: ConfigEntry[Boolean] = conf("spark.comet.testing.aggregate.partialMode.enabled") .internal() diff --git a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala index 241405c0ea5..c4d541d5e77 100644 --- a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala +++ b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala @@ -20,8 +20,10 @@ package org.apache.comet.parquet import org.apache.hadoop.conf.Configuration +import org.apache.hadoop.fs.Path import org.apache.parquet.crypto.DecryptionPropertiesFactory import org.apache.parquet.crypto.keytools.{KeyToolkit, PropertiesDrivenCryptoFactory} +import org.apache.parquet.hadoop.ParquetFileReader import org.apache.spark.sql.internal.SQLConf object CometParquetUtils { @@ -57,6 +59,33 @@ object CometParquetUtils { def ignoreMissingIds(conf: SQLConf): Boolean = conf.getConfString(IGNORE_MISSING_PARQUET_FIELD_ID, "false").toBoolean + def requiresDatetimeRebase( + paths: Seq[Path], + conf: Configuration, + datetimeMode: String, + int96Mode: String, + hasDate: Boolean, + hasTimestamp: Boolean): Boolean = { + paths.exists { path => + val reader = ParquetFileReader.open(conf, path) + try { + val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData + val version = Option(metadata.get("org.apache.spark.version")) + def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = + version.fold(mode != "CORRECTED")(v => + v < minVersion || metadata.containsKey(legacyKey)) + + (hasDate && + needsRebase(datetimeMode, "3.0.0", "org.apache.spark.legacyDateTime")) || + (hasTimestamp && + (needsRebase(datetimeMode, "3.0.0", "org.apache.spark.legacyDateTime") || + needsRebase(int96Mode, "3.1.0", "org.apache.spark.legacyINT96"))) + } finally { + reader.close() + } + } + } + /** * Checks if the given Hadoop configuration contains any unsupported encryption settings. * diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index 5e451b382ad..22525376e82 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -27,6 +27,7 @@ import java.util.concurrent.ConcurrentHashMap import scala.collection.mutable import scala.collection.mutable.ListBuffer import scala.jdk.CollectionConverters._ +import scala.util.control.NonFatal import org.apache.hadoop.conf.Configuration import org.apache.spark.internal.Logging @@ -38,6 +39,7 @@ import org.apache.spark.sql.catalyst.util.ResolveDefaultColumns.getExistenceDefa import org.apache.spark.sql.comet.{CometBatchScanExec, CometScanExec} import org.apache.spark.sql.execution.{FileSourceScanExec, InSubqueryExec, SparkPlan, SubqueryAdaptiveBroadcastExec} import org.apache.spark.sql.execution.datasources.HadoopFsRelation +import org.apache.spark.sql.execution.datasources.parquet.ParquetOptions import org.apache.spark.sql.execution.datasources.v2.BatchScanExec import org.apache.spark.sql.execution.datasources.v2.csv.CSVScan import org.apache.spark.sql.internal.SQLConf @@ -49,7 +51,8 @@ import org.apache.comet.CometSparkSessionExtensions.{isCometLoaded, isSpark35Plu import org.apache.comet.DataTypeSupport.isComplexType import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection} import org.apache.comet.objectstore.NativeConfig -import org.apache.comet.parquet.CometParquetUtils.{encryptionEnabled, isEncryptionConfigSupported} +import org.apache.comet.parquet.CometParquetUtils.{encryptionEnabled, isEncryptionConfigSupported, requiresDatetimeRebase} +import org.apache.comet.serde.SupportLevel import org.apache.comet.serde.operator.{CometIcebergNativeScan, CometNativeScan} import org.apache.comet.shims.{CometTypeShim, ShimCometStreaming, ShimFileFormat, ShimSubqueryBroadcast} @@ -295,7 +298,37 @@ case class CometScanRule(session: SparkSession) if (!isSchemaSupported(scanExec, r)) { return None } - Some(CometScanExec(scanExec, session)) + val cometScan = CometScanExec(scanExec, session) + val hasDate = SupportLevel.containsType(scanExec.requiredSchema, classOf[DateType]) + val hasTimestamp = SupportLevel.containsType( + scanExec.requiredSchema, + classOf[TimestampType], + classOf[TimestampNTZType]) + if (hasDate || hasTimestamp) { + val options = new ParquetOptions(r.options, conf) + val paths = + cometScan.selectedPartitions.iterator.flatMap(_.files.iterator.map(_.getPath)).toSeq + try { + if (requiresDatetimeRebase( + paths, + hadoopConf, + options.datetimeRebaseModeInRead, + options.int96RebaseModeInRead, + hasDate, + hasTimestamp)) { + withFallbackReason(scanExec, "Native Parquet scan does not support datetime rebasing") + return None + } + } catch { + case NonFatal(e) => + logWarning("Unable to inspect Parquet datetime rebase metadata", e) + withFallbackReason( + scanExec, + "Native Parquet scan could not verify datetime rebase metadata") + return None + } + } + Some(cometScan) } private def transformV2Scan(scanExec: BatchScanExec): SparkPlan = { diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index 688578de371..60c80c994f4 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -1733,4 +1733,27 @@ class ParquetReadV1Suite extends ParquetReadSuite with AdaptiveSparkPlanHelper { } } + test("fallback for Parquet datetime rebasing") { + withTempPath { path => + withSQLConf( + SQLConf.PARQUET_REBASE_MODE_IN_WRITE.key -> "LEGACY", + SQLConf.PARQUET_INT96_REBASE_MODE_IN_WRITE.key -> "LEGACY") { + sql(""" + |SELECT cast(s AS date) AS d, cast(s AS timestamp) AS ts + |FROM VALUES ('1000-01-01'), ('1990-01-01') AS v(s) + |""".stripMargin).write.parquet(path.toString) + } + + val df = spark.read.parquet(path.toString) + val plan = df.queryExecution.executedPlan + assert(collect(plan) { case _: CometNativeScanExec => true }.isEmpty) + checkAnswer( + df.where("d = date'1000-01-01'").selectExpr("count(*)", "min(ts)", "max(ts)"), + Row( + 1L, + java.sql.Timestamp.valueOf("1000-01-01 00:00:00"), + java.sql.Timestamp.valueOf("1000-01-01 00:00:00"))) + } + } + } From ce10dd4c18f053861eadf560b2114513a07dc9f8 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Mon, 27 Jul 2026 16:05:50 +0800 Subject: [PATCH 2/9] docs: link Spark datetime rebase policy --- .../main/scala/org/apache/comet/parquet/CometParquetUtils.scala | 2 ++ 1 file changed, 2 insertions(+) diff --git a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala index c4d541d5e77..df0491cec03 100644 --- a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala +++ b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala @@ -71,6 +71,8 @@ object CometParquetUtils { try { val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData val version = Option(metadata.get("org.apache.spark.version")) + // Match Spark's metadata precedence and version thresholds: + // https://github.com/apache/spark/blob/v4.1.2/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/DataSourceUtils.scala#L130-L176 def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = version.fold(mode != "CORRECTED")(v => v < minVersion || metadata.containsKey(legacyKey)) From b8bba9a14bc385db873593bffb856c32dc185dd6 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Mon, 27 Jul 2026 16:10:46 +0800 Subject: [PATCH 3/9] upd reference --- .../scala/org/apache/comet/parquet/CometParquetUtils.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala index df0491cec03..e7b086a2da1 100644 --- a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala +++ b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala @@ -71,8 +71,8 @@ object CometParquetUtils { try { val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData val version = Option(metadata.get("org.apache.spark.version")) - // Match Spark's metadata precedence and version thresholds: - // https://github.com/apache/spark/blob/v4.1.2/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/DataSourceUtils.scala#L130-L176 + // Copy from Spark's rebase spec + // https://github.com/apache/spark/blob/v4.1.2/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/DataSourceUtils.scala#L130-L177 def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = version.fold(mode != "CORRECTED")(v => v < minVersion || metadata.containsKey(legacyKey)) From aa4cda1ccf70cd0846627ac236efd837c81684a8 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Mon, 27 Jul 2026 20:03:23 +0800 Subject: [PATCH 4/9] fix: allow long Spark source reference --- .../main/scala/org/apache/comet/parquet/CometParquetUtils.scala | 2 ++ 1 file changed, 2 insertions(+) diff --git a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala index e7b086a2da1..3e23809028f 100644 --- a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala +++ b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala @@ -72,7 +72,9 @@ object CometParquetUtils { val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData val version = Option(metadata.get("org.apache.spark.version")) // Copy from Spark's rebase spec + // scalastyle:off line.size.limit // https://github.com/apache/spark/blob/v4.1.2/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/DataSourceUtils.scala#L130-L177 + // scalastyle:on line.size.limit def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = version.fold(mode != "CORRECTED")(v => v < minVersion || metadata.containsKey(legacyKey)) From ac2a5af0a06502353f597fd1bccddac4202504a7 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Mon, 27 Jul 2026 20:05:44 +0800 Subject: [PATCH 5/9] remove comments --- .../scala/org/apache/comet/parquet/CometParquetUtils.scala | 4 ---- 1 file changed, 4 deletions(-) diff --git a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala index 3e23809028f..c4d541d5e77 100644 --- a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala +++ b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala @@ -71,10 +71,6 @@ object CometParquetUtils { try { val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData val version = Option(metadata.get("org.apache.spark.version")) - // Copy from Spark's rebase spec - // scalastyle:off line.size.limit - // https://github.com/apache/spark/blob/v4.1.2/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/DataSourceUtils.scala#L130-L177 - // scalastyle:on line.size.limit def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = version.fold(mode != "CORRECTED")(v => v < minVersion || metadata.containsKey(legacyKey)) From f8ba5cddf8b325fb3224483a025c14f6fc729b1e Mon Sep 17 00:00:00 2001 From: peterxcli Date: Mon, 27 Jul 2026 21:06:08 +0800 Subject: [PATCH 6/9] fix: preserve corrected Parquet scan provenance --- native/core/src/execution/operators/parquet_writer.rs | 6 +++++- .../org/apache/comet/parquet/CometParquetUtils.scala | 4 +++- .../org/apache/comet/parquet/ParquetReadSuite.scala | 9 ++------- .../test/scala/org/apache/spark/sql/CometTestBase.scala | 2 ++ 4 files changed, 12 insertions(+), 9 deletions(-) diff --git a/native/core/src/execution/operators/parquet_writer.rs b/native/core/src/execution/operators/parquet_writer.rs index dbbee713ae1..3a5bc0b450e 100644 --- a/native/core/src/execution/operators/parquet_writer.rs +++ b/native/core/src/execution/operators/parquet_writer.rs @@ -52,7 +52,7 @@ use futures::TryStreamExt; use parquet::{ arrow::ArrowWriter, basic::{Compression, GzipLevel, ZstdLevel}, - file::properties::WriterProperties, + file::{metadata::KeyValue, properties::WriterProperties}, }; use url::Url; @@ -506,6 +506,10 @@ impl ExecutionPlan for ParquetWriterExec { // Configure writer properties let props = WriterProperties::builder() .set_compression(compression) + .set_key_value_metadata(Some(vec![KeyValue::new( + "org.apache.comet.datetimeRebaseMode".to_string(), + Some("CORRECTED".to_string()), + )])) .build(); let object_store_options = self.object_store_options.clone(); diff --git a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala index c4d541d5e77..e59c27580d3 100644 --- a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala +++ b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala @@ -71,8 +71,10 @@ object CometParquetUtils { try { val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData val version = Option(metadata.get("org.apache.spark.version")) + val cometCorrected = + Option(metadata.get("org.apache.comet.datetimeRebaseMode")).contains("CORRECTED") def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = - version.fold(mode != "CORRECTED")(v => + version.fold(!cometCorrected && mode != "CORRECTED")(v => v < minVersion || metadata.containsKey(legacyKey)) (hasDate && diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index 60c80c994f4..69804fca1ba 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -1711,20 +1711,15 @@ class ParquetReadV1Suite extends ParquetReadSuite with AdaptiveSparkPlanHelper { } } - test("reading ancient dates before 1582") { - // Verify that legacy dates (before 1582-10-15) are read without error. - // Comet does not support datetime rebasing, so these dates are read as if they were - // written using the Proleptic Gregorian calendar (no rebase, no exception). + test("fallback when reading ancient dates before 1582") { val file = getResourceParquetFilePath("test-data/before_1582_date_v3_2_0.snappy.parquet") val df = spark.read.parquet(file) - // Verify Comet scan is in the plan val plan = df.queryExecution.executedPlan - checkCometOperators(plan) + assert(collect(plan) { case _: CometNativeScanExec => true }.isEmpty) - // Verify all 8 rows are read and contain dates before 1582 val rows = df.collect() assert(rows.length == 8, s"Expected 8 rows, got ${rows.length}") rows.foreach { row => diff --git a/spark/src/test/scala/org/apache/spark/sql/CometTestBase.scala b/spark/src/test/scala/org/apache/spark/sql/CometTestBase.scala index 9a59c9c2ebf..db96535be78 100644 --- a/spark/src/test/scala/org/apache/spark/sql/CometTestBase.scala +++ b/spark/src/test/scala/org/apache/spark/sql/CometTestBase.scala @@ -677,6 +677,8 @@ abstract class CometTestBase .withPageSize(pageSize) .withDictionaryPageSize(dictionaryPageSize) .withPageRowCountLimit(pageRowCountLimit) + .withExtraMetaData( + java.util.Collections.singletonMap("org.apache.spark.version", SPARK_VERSION)) .withConf(hadoopConf) .build() } From f74c94cad6eea95b9fc6f5943ba1683493bce871 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Tue, 4 Aug 2026 02:32:56 +0800 Subject: [PATCH 7/9] fix: reduce datetime rebase fallback planning cost --- .../user-guide/latest/compatibility/scans.md | 11 --- .../comet/parquet/CometParquetUtils.scala | 31 ------ .../apache/comet/rules/CometScanRule.scala | 14 +-- .../spark/sql/comet/CometScanUtils.scala | 75 ++++++++++++++ .../comet/parquet/ParquetReadSuite.scala | 97 ++++++++++++++++--- 5 files changed, 163 insertions(+), 65 deletions(-) diff --git a/docs/source/user-guide/latest/compatibility/scans.md b/docs/source/user-guide/latest/compatibility/scans.md index ce3a610ec56..3e1dc104639 100644 --- a/docs/source/user-guide/latest/compatibility/scans.md +++ b/docs/source/user-guide/latest/compatibility/scans.md @@ -54,17 +54,6 @@ The following features are not supported and cause Comet to fall back to Spark: `spark.comet.scan.allowDisabledParquetVectorizedReader=true` to opt in to running the Comet Parquet scan regardless. -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 with `spark.sql.parquet.datetimeRebaseModeInWrite=LEGACY` (which is Spark's default for - data written before Spark 3.0, using the hybrid Julian/Gregorian calendar), Comet reads them as - if they were written using the Proleptic Gregorian calendar. This produces silently-wrong - values for dates before October 15, 1582 in both projections and predicates. Comet also - ignores `spark.sql.parquet.datetimeRebaseModeInRead` and the file-level - `org.apache.spark.legacyDateTime` metadata that would tell it to rebase. Tracked by - [#5010](https://github.com/apache/datafusion-comet/issues/5010). - The following limitations raise an error at scan time rather than falling back to Spark: - Invalid UTF-8 bytes in `STRING` columns. Spark permits arbitrary byte sequences in a `STRING` diff --git a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala index e59c27580d3..241405c0ea5 100644 --- a/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala +++ b/spark/src/main/scala/org/apache/comet/parquet/CometParquetUtils.scala @@ -20,10 +20,8 @@ package org.apache.comet.parquet import org.apache.hadoop.conf.Configuration -import org.apache.hadoop.fs.Path import org.apache.parquet.crypto.DecryptionPropertiesFactory import org.apache.parquet.crypto.keytools.{KeyToolkit, PropertiesDrivenCryptoFactory} -import org.apache.parquet.hadoop.ParquetFileReader import org.apache.spark.sql.internal.SQLConf object CometParquetUtils { @@ -59,35 +57,6 @@ object CometParquetUtils { def ignoreMissingIds(conf: SQLConf): Boolean = conf.getConfString(IGNORE_MISSING_PARQUET_FIELD_ID, "false").toBoolean - def requiresDatetimeRebase( - paths: Seq[Path], - conf: Configuration, - datetimeMode: String, - int96Mode: String, - hasDate: Boolean, - hasTimestamp: Boolean): Boolean = { - paths.exists { path => - val reader = ParquetFileReader.open(conf, path) - try { - val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData - val version = Option(metadata.get("org.apache.spark.version")) - val cometCorrected = - Option(metadata.get("org.apache.comet.datetimeRebaseMode")).contains("CORRECTED") - def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = - version.fold(!cometCorrected && mode != "CORRECTED")(v => - v < minVersion || metadata.containsKey(legacyKey)) - - (hasDate && - needsRebase(datetimeMode, "3.0.0", "org.apache.spark.legacyDateTime")) || - (hasTimestamp && - (needsRebase(datetimeMode, "3.0.0", "org.apache.spark.legacyDateTime") || - needsRebase(int96Mode, "3.1.0", "org.apache.spark.legacyINT96"))) - } finally { - reader.close() - } - } - } - /** * Checks if the given Hadoop configuration contains any unsupported encryption settings. * diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index d76b025bf3c..1ddaf32ddf4 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -36,7 +36,7 @@ import org.apache.spark.sql.catalyst.expressions.{Attribute, DynamicPruningExpre import org.apache.spark.sql.catalyst.rules.Rule import org.apache.spark.sql.catalyst.util.{sideBySide, ArrayBasedMapData, GenericArrayData, MetadataColumnHelper} import org.apache.spark.sql.catalyst.util.ResolveDefaultColumns.getExistenceDefaultValues -import org.apache.spark.sql.comet.{CometBatchScanExec, CometScanExec} +import org.apache.spark.sql.comet.{CometBatchScanExec, CometScanExec, CometScanUtils} import org.apache.spark.sql.execution.{FileSourceScanExec, InSubqueryExec, SparkPlan, SubqueryAdaptiveBroadcastExec} import org.apache.spark.sql.execution.datasources.HadoopFsRelation import org.apache.spark.sql.execution.datasources.parquet.ParquetOptions @@ -51,7 +51,7 @@ import org.apache.comet.CometSparkSessionExtensions.{isCometLoaded, isSpark35Plu import org.apache.comet.DataTypeSupport.isComplexType import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection} import org.apache.comet.objectstore.NativeConfig -import org.apache.comet.parquet.CometParquetUtils.{encryptionEnabled, isEncryptionConfigSupported, requiresDatetimeRebase} +import org.apache.comet.parquet.CometParquetUtils.{encryptionEnabled, isEncryptionConfigSupported} import org.apache.comet.serde.SupportLevel import org.apache.comet.serde.operator.{CometIcebergNativeScan, CometNativeScan} import org.apache.comet.shims.{CometTypeShim, ShimCometStreaming, ShimFileFormat, ShimSubqueryBroadcast} @@ -295,16 +295,16 @@ case class CometScanRule(session: SparkSession) } val cometScan = CometScanExec(scanExec, session) val hasDate = SupportLevel.containsType(scanExec.requiredSchema, classOf[DateType]) - val hasTimestamp = SupportLevel.containsType( - scanExec.requiredSchema, - classOf[TimestampType], - classOf[TimestampNTZType]) + val hasTimestamp = + SupportLevel.containsType(scanExec.requiredSchema, classOf[TimestampType]) || + (COMET_SCHEMA_EVOLUTION_ENABLED && + SupportLevel.containsType(scanExec.requiredSchema, classOf[TimestampNTZType])) if (hasDate || hasTimestamp) { val options = new ParquetOptions(r.options, conf) val paths = cometScan.selectedPartitions.iterator.flatMap(_.files.iterator.map(_.getPath)).toSeq try { - if (requiresDatetimeRebase( + if (CometScanUtils.requiresDatetimeRebase( paths, hadoopConf, options.datetimeRebaseModeInRead, diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala index 19a2b53af11..ae3be5ffc25 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala @@ -19,11 +19,86 @@ package org.apache.spark.sql.comet +import java.util.concurrent.{Callable, ExecutorCompletionService} + +import org.apache.hadoop.conf.Configuration +import org.apache.hadoop.fs.Path +import org.apache.parquet.HadoopReadOptions +import org.apache.parquet.format.converter.ParquetMetadataConverter.SKIP_ROW_GROUPS +import org.apache.parquet.hadoop.ParquetFileReader +import org.apache.parquet.hadoop.util.HadoopInputFile import org.apache.spark.sql.catalyst.expressions.{DynamicPruningExpression, Expression, Literal} import org.apache.spark.sql.execution.{InSubqueryExec, SubqueryAdaptiveBroadcastExec} +import org.apache.spark.util.ThreadUtils object CometScanUtils { + def requiresDatetimeRebase( + paths: Seq[Path], + conf: Configuration, + datetimeMode: String, + int96Mode: String, + hasDate: Boolean, + hasTimestamp: Boolean): Boolean = { + val parallelism = 8 + val pool = ThreadUtils.newDaemonFixedThreadPool(parallelism, "checkingParquetDatetimeRebase") + val completion = new ExecutorCompletionService[Boolean](pool) + val remaining = paths.iterator + var inFlight = 0 + + // Spark's Parquet footer reader uses ThreadUtils.parmap, which submits every input eagerly. + // Keep only `parallelism` reads in flight so finding one legacy footer stops further reads. + + def submitNext(): Unit = { + val path = remaining.next() + completion.submit(new Callable[Boolean] { + override def call(): Boolean = { + val inputFile = HadoopInputFile.fromPath(path, conf) + val readOptions = HadoopReadOptions + .builder(conf, path) + .withMetadataFilter(SKIP_ROW_GROUPS) + .build() + val reader = ParquetFileReader.open(inputFile, readOptions) + try { + val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData + val version = Option(metadata.get("org.apache.spark.version")) + val cometCorrected = + Option(metadata.get("org.apache.comet.datetimeRebaseMode")).contains("CORRECTED") + def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = + version.fold(!cometCorrected && mode != "CORRECTED")(v => + v < minVersion || metadata.containsKey(legacyKey)) + + (hasDate && + needsRebase(datetimeMode, "3.0.0", "org.apache.spark.legacyDateTime")) || + (hasTimestamp && + (needsRebase(datetimeMode, "3.0.0", "org.apache.spark.legacyDateTime") || + needsRebase(int96Mode, "3.1.0", "org.apache.spark.legacyINT96"))) + } finally { + reader.close() + } + } + }) + inFlight += 1 + } + + try { + while (inFlight < parallelism && remaining.hasNext) { + submitNext() + } + var requiresRebase = false + while (!requiresRebase && inFlight > 0) { + requiresRebase = completion.take().get() + inFlight -= 1 + if (!requiresRebase && remaining.hasNext) { + submitNext() + } + } + requiresRebase + } finally { + pool.shutdownNow() + } + } + /** * Filters unused DynamicPruningExpression expressions - one which has been replaced with * DynamicPruningExpression(Literal.TrueLiteral) during Physical Planning diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index 9b2897950dd..35bb56faeb7 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -35,7 +35,7 @@ import org.apache.parquet.schema.MessageTypeParser import org.apache.spark.SparkException import org.apache.spark.sql.{CometTestBase, DataFrame, Row} import org.apache.spark.sql.catalyst.util.DateTimeUtils -import org.apache.spark.sql.comet.{CometNativeScanExec, CometScanExec} +import org.apache.spark.sql.comet.{CometNativeScanExec, CometScanExec, CometScanUtils} import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper import org.apache.spark.sql.execution.datasources.parquet.ParquetUtils import org.apache.spark.sql.internal.SQLConf @@ -1711,37 +1711,83 @@ class ParquetReadV1Suite extends ParquetReadSuite with AdaptiveSparkPlanHelper { } } - test("fallback when reading ancient dates before 1582") { - val file = - getResourceParquetFilePath("test-data/before_1582_date_v3_2_0.snappy.parquet") + test("fallback for ancient datetime fixtures") { + val encodings = Seq( + "date", + "timestamp_micros", + "timestamp_millis", + "timestamp_int96_plain", + "timestamp_int96_dict") + val versions = Seq("v2_4_5", "v2_4_6", "v3_2_0") - val df = spark.read.parquet(file) - - val plan = df.queryExecution.executedPlan - assert(collect(plan) { case _: CometNativeScanExec => true }.isEmpty) - - val rows = df.collect() - assert(rows.length == 8, s"Expected 8 rows, got ${rows.length}") - rows.foreach { row => - val date = row.getDate(0) - assert(date.toLocalDate.getYear < 1582, s"Expected date before 1582, got $date") + withSQLConf( + "spark.sql.parquet.datetimeRebaseModeInRead" -> "LEGACY", + "spark.sql.parquet.int96RebaseModeInRead" -> "LEGACY") { + for (encoding <- encodings; version <- versions) { + val name = s"before_1582_${encoding}_$version.snappy.parquet" + withClue(s"$name: ") { + val (_, cometPlan) = + checkSparkAnswer(spark.read.parquet(getResourceParquetFilePath(s"test-data/$name"))) + assert(collect(cometPlan) { case _: CometNativeScanExec => true }.isEmpty) + } + } } } test("fallback for Parquet datetime rebasing") { withTempPath { path => + val correctedPath = new File(path, "corrected") + val legacyPath = new File(path, "legacy") + withSQLConf( + SQLConf.PARQUET_REBASE_MODE_IN_WRITE.key -> "CORRECTED", + SQLConf.PARQUET_INT96_REBASE_MODE_IN_WRITE.key -> "CORRECTED") { + sql(""" + |SELECT cast(s AS date) AS d, cast(s AS timestamp) AS ts + |FROM VALUES ('2000-01-01') AS v(s) + |""".stripMargin).coalesce(1).write.parquet(correctedPath.toString) + } withSQLConf( SQLConf.PARQUET_REBASE_MODE_IN_WRITE.key -> "LEGACY", SQLConf.PARQUET_INT96_REBASE_MODE_IN_WRITE.key -> "LEGACY") { sql(""" |SELECT cast(s AS date) AS d, cast(s AS timestamp) AS ts |FROM VALUES ('1000-01-01'), ('1990-01-01') AS v(s) - |""".stripMargin).write.parquet(path.toString) + |""".stripMargin).coalesce(1).write.parquet(legacyPath.toString) } - val df = spark.read.parquet(path.toString) + val hadoopConf = spark.sessionState.newHadoopConf() + def parquetFile(dir: File): Path = { + val dirPath = new Path(dir.toString) + dirPath + .getFileSystem(hadoopConf) + .listStatus(dirPath) + .iterator + .map(_.getPath) + .find(_.getName.endsWith(".parquet")) + .get + } + def requiresRebase(paths: Seq[Path]): Boolean = + CometScanUtils.requiresDatetimeRebase( + paths, + hadoopConf, + "CORRECTED", + "CORRECTED", + hasDate = true, + hasTimestamp = true) + + val correctedFile = parquetFile(correctedPath) + val legacyFile = parquetFile(legacyPath) + assert(requiresRebase(Seq(correctedFile, legacyFile))) + assert( + requiresRebase( + Seq.fill(8)(legacyFile) :+ new Path(path.toString, "must-not-be-read.parquet"))) + + val df = spark.read.parquet(correctedPath.toString, legacyPath.toString) val plan = df.queryExecution.executedPlan assert(collect(plan) { case _: CometNativeScanExec => true }.isEmpty) + checkAnswer( + df.selectExpr("count(*)", "min(d)", "max(d)"), + Row(3L, java.sql.Date.valueOf("1000-01-01"), java.sql.Date.valueOf("2000-01-01"))) checkAnswer( df.where("d = date'1000-01-01'").selectExpr("count(*)", "min(ts)", "max(ts)"), Row( @@ -1751,4 +1797,23 @@ class ParquetReadV1Suite extends ParquetReadSuite with AdaptiveSparkPlanHelper { } } + test("timestamp_ntz falls back for legacy datetime metadata on Spark 4+") { + withTempPath { path => + withSQLConf(SQLConf.PARQUET_REBASE_MODE_IN_WRITE.key -> "LEGACY") { + sql(""" + |SELECT cast('1000-01-01 00:00:00' AS timestamp_ntz) AS ts_ntz, + | cast('1000-01-01' AS date) AS legacy_date + |""".stripMargin).coalesce(1).write.parquet(path.toString) + } + + val (_, cometPlan) = checkSparkAnswer(spark.read.parquet(path.toString).select("ts_ntz")) + val nativeScans = collect(cometPlan) { case _: CometNativeScanExec => true } + if (CometConf.COMET_SCHEMA_EVOLUTION_ENABLED) { + assert(nativeScans.isEmpty) + } else { + assert(nativeScans.nonEmpty) + } + } + } + } From 2b166fe3a9d916bd942664f5c5364b3282c2c85f Mon Sep 17 00:00:00 2001 From: peterxcli Date: Tue, 4 Aug 2026 15:54:54 +0800 Subject: [PATCH 8/9] fix: write Spark version to Parquet metadata --- .../src/execution/operators/parquet_writer.rs | 15 ++++++++++++--- native/core/src/execution/planner.rs | 1 + native/proto/src/proto/operator.proto | 2 ++ .../serde/operator/CometDataWritingCommand.scala | 3 ++- .../apache/spark/sql/comet/CometScanUtils.scala | 4 +--- .../comet/parquet/CometParquetWriterSuite.scala | 11 +++++++++++ 6 files changed, 29 insertions(+), 7 deletions(-) diff --git a/native/core/src/execution/operators/parquet_writer.rs b/native/core/src/execution/operators/parquet_writer.rs index 3a5bc0b450e..f7e37cd0667 100644 --- a/native/core/src/execution/operators/parquet_writer.rs +++ b/native/core/src/execution/operators/parquet_writer.rs @@ -232,6 +232,8 @@ pub struct ParquetWriterExec { partition_id: i32, /// Column names to use in the output Parquet file column_names: Vec, + /// Runtime Spark version to record in the Parquet metadata + spark_version: String, /// Object store configuration options object_store_options: HashMap, /// Metrics @@ -252,6 +254,7 @@ impl ParquetWriterExec { compression: ParquetCompression, partition_id: i32, column_names: Vec, + spark_version: String, object_store_options: HashMap, ) -> Result { // Preserve the input's partitioning so each partition writes its own file @@ -273,6 +276,7 @@ impl ParquetWriterExec { compression, partition_id, column_names, + spark_version, object_store_options, metrics: ExecutionPlanMetricsSet::new(), cache, @@ -453,6 +457,7 @@ impl ExecutionPlan for ParquetWriterExec { self.compression.clone(), self.partition_id, self.column_names.clone(), + self.spark_version.clone(), self.object_store_options.clone(), )?)), _ => Err(DataFusionError::Internal( @@ -506,9 +511,12 @@ impl ExecutionPlan for ParquetWriterExec { // Configure writer properties let props = WriterProperties::builder() .set_compression(compression) + // Spark identifies corrected datetime files by its writer version and the absence of + // legacy markers. Comet always writes corrected values, so use the same metadata: + // https://github.com/apache/spark/blob/v4.2.0/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetWriteSupport.scala#L126-L146 .set_key_value_metadata(Some(vec![KeyValue::new( - "org.apache.comet.datetimeRebaseMode".to_string(), - Some("CORRECTED".to_string()), + "org.apache.spark.version".to_string(), + Some(self.spark_version.clone()), )])) .build(); @@ -868,7 +876,8 @@ mod tests { ParquetCompression::None, 0, // partition_id column_names, - HashMap::new(), // object_store_options + "4.2.0".to_string(), // spark_version + HashMap::new(), // object_store_options )?; // Create a session context and execute the plan diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index f20dadf7f3c..2fb15b1c93c 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -1828,6 +1828,7 @@ impl PhysicalPlanner { codec, self.partition, writer.column_names.clone(), + writer.spark_version.clone(), object_store_options, )?); diff --git a/native/proto/src/proto/operator.proto b/native/proto/src/proto/operator.proto index ced87262f32..512e9f4e2e5 100644 --- a/native/proto/src/proto/operator.proto +++ b/native/proto/src/proto/operator.proto @@ -475,6 +475,8 @@ message ParquetWriter { // configuration value "spark.hadoop.fs.s3a.access.key" will be stored as "fs.s3a.access.key" in // the map. map object_store_options = 8; + // Runtime Spark version written to Parquet metadata for Spark reader compatibility. + string spark_version = 9; } enum AggregateMode { diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala index 8157f286825..f37fb60ca8c 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala @@ -25,7 +25,7 @@ import java.util.Locale import scala.jdk.CollectionConverters._ import org.apache.parquet.hadoop.ParquetOutputFormat -import org.apache.spark.SparkException +import org.apache.spark.{SPARK_VERSION_SHORT, SparkException} import org.apache.spark.sql.comet.{CometNativeExec, CometNativeWriteExec} import org.apache.spark.sql.execution.command.DataWritingCommandExec import org.apache.spark.sql.execution.datasources.{InsertIntoHadoopFsRelationCommand, WriteFilesExec} @@ -134,6 +134,7 @@ object CometDataWritingCommand extends CometOperatorSerde[DataWritingCommandExec .newBuilder() .setOutputPath(outputPath) .setCompression(codec) + .setSparkVersion(SPARK_VERSION_SHORT) .addAllColumnNames(cmd.query.output.map(_.name).asJava) // Note: work_dir, job_id, and task_attempt_id will be set at execution time // in CometNativeWriteExec, as they depend on the Spark task context diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala index ae3be5ffc25..8f8e619fbf3 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala @@ -62,10 +62,8 @@ object CometScanUtils { try { val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData val version = Option(metadata.get("org.apache.spark.version")) - val cometCorrected = - Option(metadata.get("org.apache.comet.datetimeRebaseMode")).contains("CORRECTED") def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = - version.fold(!cometCorrected && mode != "CORRECTED")(v => + version.fold(mode != "CORRECTED")(v => v < minVersion || metadata.containsKey(legacyKey)) (hasDate && diff --git a/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala index a1ae1af1d1c..ae7766bfbac 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala @@ -28,6 +28,7 @@ import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.parquet.hadoop.ParquetFileReader import org.apache.parquet.hadoop.metadata.CompressionCodecName import org.apache.parquet.hadoop.util.HadoopInputFile +import org.apache.spark.SPARK_VERSION_SHORT import org.apache.spark.sql.{AnalysisException, CometTestBase, DataFrame, Row, SaveMode} import org.apache.spark.sql.comet.{CometBatchScanExec, CometNativeScanExec, CometNativeWriteExec, CometScanExec} import org.apache.spark.sql.execution.{FileSourceScanExec, QueryExecution, SparkPlan} @@ -812,6 +813,16 @@ class CometParquetWriterSuite extends CometTestBase { // With 1000 rows and default parallelism, we should get multiple partitions assert(partFiles.length > 1, "Expected multiple part files to be created") + val conf = spark.sparkContext.hadoopConfiguration + partFiles.foreach { partFile => + val inputFile = HadoopInputFile.fromPath(new Path(partFile.getAbsolutePath), conf) + Using.resource(ParquetFileReader.open(inputFile)) { reader => + val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData + assert(metadata.get("org.apache.spark.version") == SPARK_VERSION_SHORT) + assert(!metadata.containsKey("org.apache.comet.datetimeRebaseMode")) + } + } + // read with and without Comet and compare val sparkRows = readSparkRows(outputPath) val cometRows = readCometRows(outputPath) From 2323bae3a9d01b69a34067a6f884b7e29aea5d88 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Fri, 28 Aug 2026 01:09:15 +0800 Subject: [PATCH 9/9] fix: address review feedback on datetime rebase fallback - Cache per-file footer facts (keyed by path, length, modification time) so AQE re-planning and repeated queries do not re-read footers - Add spark.comet.scan.parquet.checkDatetimeRebase config to opt out of the planning-time footer check - Fall back to Spark for native Parquet writes when the write rebase mode is LEGACY and the schema contains dates or timestamps - Gate the TimestampNTZ rebase check on COMET_ALLOW_TIMESTAMP_LTZ_AS_NTZ instead of COMET_SCHEMA_EVOLUTION_ENABLED, with an explanation - Document that the version comparison mirrors Spark's DataSourceUtils.datetimeRebaseSpec Co-Authored-By: Claude Fable 5 --- .../user-guide/latest/compatibility/scans.md | 6 +- .../scala/org/apache/comet/CometConf.scala | 16 +++ .../apache/comet/rules/CometScanRule.scala | 18 ++- .../operator/CometDataWritingCommand.scala | 29 +++++ .../spark/sql/comet/CometScanUtils.scala | 119 +++++++++++++----- .../parquet/CometParquetWriterSuite.scala | 37 ++++++ .../comet/parquet/ParquetReadSuite.scala | 47 +++++-- 7 files changed, 229 insertions(+), 43 deletions(-) diff --git a/docs/source/user-guide/latest/compatibility/scans.md b/docs/source/user-guide/latest/compatibility/scans.md index 3e1dc104639..72c69ff9440 100644 --- a/docs/source/user-guide/latest/compatibility/scans.md +++ b/docs/source/user-guide/latest/compatibility/scans.md @@ -47,7 +47,11 @@ The following features are not supported and cause Comet to fall back to Spark: `spark.sql.parquet.datetimeRebaseModeInRead` / `spark.sql.parquet.int96RebaseModeInRead` requires legacy-calendar handling. This includes files written by any Spark version with a corresponding legacy rebase mode, not only files - written before Spark 3.0. See [#5010](https://github.com/apache/datafusion-comet/issues/5010). + written before Spark 3.0. Detecting this requires reading each input file's footer on the + driver during planning (results are cached per file); users whose data is known to be free + of legacy-calendar values can skip the check by setting + `spark.comet.scan.parquet.checkDatetimeRebase=false`. + See [#5010](https://github.com/apache/datafusion-comet/issues/5010). - `spark.sql.parquet.enableVectorizedReader=false`. Disabling the vectorized reader opts into Spark's parquet-mr semantics (silent overflow, null-on-narrowing), which Comet's native reader does not replicate. By default Comet falls back to Spark in this case. Set diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index ae297151932..bc8443392e9 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -771,6 +771,22 @@ object CometConf extends ShimCometConf { .booleanConf .createWithDefault(true) + val COMET_SCAN_PARQUET_CHECK_DATETIME_REBASE: ConfigEntry[Boolean] = + conf("spark.comet.scan.parquet.checkDatetimeRebase") + .category(CATEGORY_SCAN) + .doc( + "Whether to inspect Parquet footer metadata during planning to detect files whose " + + "dates/timestamps may require legacy (hybrid Julian/Gregorian) datetime rebasing, " + + "and fall back to Spark for those scans. The check reads each input file's footer " + + "on the driver the first time the file is planned; results are cached per file. " + + "Disable only when all input files are known to contain datetime values written " + + "with the proleptic Gregorian calendar (for example, written by Spark 3.x or later " + + "with corrected rebase modes). When disabled, Comet reads legacy files without " + + "rebasing, which produces results that differ from Spark for dates and timestamps " + + s"before 1582-10-15. $COMPAT_GUIDE.") + .booleanConf + .createWithDefault(true) + val COMET_SCAN_ALLOW_DISABLED_PARQUET_VECTORIZED_READER: ConfigEntry[Boolean] = conf("spark.comet.scan.allowDisabledParquetVectorizedReader") .category(CATEGORY_SCAN) diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index 1ddaf32ddf4..30d293eb946 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -295,17 +295,25 @@ case class CometScanRule(session: SparkSession) } val cometScan = CometScanExec(scanExec, session) val hasDate = SupportLevel.containsType(scanExec.requiredSchema, classOf[DateType]) + // TIMESTAMP_NTZ values are never rebased by Spark, on write or on read (Spark's + // ParquetVectorUpdaterFactory: "TIMESTAMP_NTZ is a new data type and has no legacy files + // that need to do rebase"). The rebase question arises for a requested NTZ column only + // when the underlying Parquet column is a TIMESTAMP (LTZ or INT96) that may carry + // legacy-calendar values, and Comet permits that read only when + // COMET_ALLOW_TIMESTAMP_LTZ_AS_NTZ is true (Spark 4.x, SPARK-47447). val hasTimestamp = SupportLevel.containsType(scanExec.requiredSchema, classOf[TimestampType]) || - (COMET_SCHEMA_EVOLUTION_ENABLED && + (COMET_ALLOW_TIMESTAMP_LTZ_AS_NTZ && SupportLevel.containsType(scanExec.requiredSchema, classOf[TimestampNTZType])) - if (hasDate || hasTimestamp) { + if ((hasDate || hasTimestamp) && COMET_SCAN_PARQUET_CHECK_DATETIME_REBASE.get()) { val options = new ParquetOptions(r.options, conf) - val paths = - cometScan.selectedPartitions.iterator.flatMap(_.files.iterator.map(_.getPath)).toSeq + val files = cometScan.selectedPartitions.iterator + .flatMap(_.files.iterator.map(f => + CometScanUtils.ParquetFileInfo(f.getPath, f.getLen, f.getModificationTime))) + .toSeq try { if (CometScanUtils.requiresDatetimeRebase( - paths, + files, hadoopConf, options.datetimeRebaseModeInRead, options.int96RebaseModeInRead, diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala index f37fb60ca8c..b4275d79d72 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala @@ -31,6 +31,7 @@ import org.apache.spark.sql.execution.command.DataWritingCommandExec import org.apache.spark.sql.execution.datasources.{InsertIntoHadoopFsRelationCommand, WriteFilesExec} import org.apache.spark.sql.execution.datasources.parquet.ParquetFileFormat import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.types.{DateType, TimestampType} import org.apache.comet.{CometConf, ConfigEntry} import org.apache.comet.CometSparkSessionExtensions.withFallbackReason @@ -78,6 +79,34 @@ object CometDataWritingCommand extends CometOperatorSerde[DataWritingCommandExec return Unsupported(Some(s"Unsupported compression codec: $codec")) } + // The native writer always writes proleptic Gregorian (corrected) datetime values + // and stamps `org.apache.spark.version` with no legacy markers. Honoring a LEGACY + // write rebase mode would require rebasing the values and stamping + // `org.apache.spark.legacyDateTime` / `org.apache.spark.legacyINT96`, so fall back + // to Spark rather than silently ignoring the requested mode and letting readers + // trust a "corrected" marker over legacy-intent data. TIMESTAMP_NTZ is exempt + // because Spark never rebases NTZ values on write. + val hasDate = cmd.query.output.exists(a => + SupportLevel.containsType(a.dataType, classOf[DateType])) + val hasTimestamp = cmd.query.output.exists(a => + SupportLevel.containsType(a.dataType, classOf[TimestampType])) + // Both write rebase mode configs default to EXCEPTION in all supported Spark + // versions. + def isLegacyWriteMode(key: String): Boolean = + SQLConf.get.getConfString(key, "EXCEPTION").toUpperCase(Locale.ROOT) == "LEGACY" + val legacyModeKeys = + ((if (hasDate || hasTimestamp) Seq(SQLConf.PARQUET_REBASE_MODE_IN_WRITE.key) + else Seq.empty) ++ + (if (hasTimestamp) Seq(SQLConf.PARQUET_INT96_REBASE_MODE_IN_WRITE.key) + else Seq.empty)).filter(isLegacyWriteMode) + if (legacyModeKeys.nonEmpty) { + return Unsupported( + Some( + "Native Parquet write always writes corrected (proleptic Gregorian) " + + "datetime values and does not support LEGACY rebase mode " + + s"(${legacyModeKeys.mkString(", ")})")) + } + Incompatible(Some("Parquet write support is highly experimental")) case _ => Unsupported(Some("Only Parquet writes are supported")) diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala index 8f8e619fbf3..bc495accd2c 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometScanUtils.scala @@ -21,6 +21,8 @@ package org.apache.spark.sql.comet import java.util.concurrent.{Callable, ExecutorCompletionService} +import scala.collection.mutable.ListBuffer + import org.apache.hadoop.conf.Configuration import org.apache.hadoop.fs.Path import org.apache.parquet.HadoopReadOptions @@ -33,48 +35,105 @@ import org.apache.spark.util.ThreadUtils object CometScanUtils { + /** Identity of one Parquet file for the datetime rebase check. */ + case class ParquetFileInfo(path: Path, length: Long, modificationTime: Long) + + /** Datetime-relevant footer metadata of one Parquet file, independent of any read mode. */ + private case class DatetimeFooterFacts( + sparkVersion: Option[String], + hasLegacyDateTime: Boolean, + hasLegacyInt96: Boolean) + + private type FooterCacheKey = (String, Long, Long) + + private val footerFactsCacheMaxSize = 32 * 1024 + + // Bounded LRU cache of per-file footer facts, keyed by (path, length, modificationTime) so a + // rewritten file is re-read. The cached facts are independent of the read modes and requested + // types, so one entry answers the rebase question for any query. This keeps AQE re-planning + // and repeated queries over the same files from re-reading footers on the driver. + private val footerFactsCache = + java.util.Collections.synchronizedMap( + new java.util.LinkedHashMap[FooterCacheKey, DatetimeFooterFacts](64, 0.75f, true) { + override def removeEldestEntry( + eldest: java.util.Map.Entry[FooterCacheKey, DatetimeFooterFacts]): Boolean = + size() > footerFactsCacheMaxSize + }) + def requiresDatetimeRebase( - paths: Seq[Path], + files: Seq[ParquetFileInfo], conf: Configuration, datetimeMode: String, int96Mode: String, hasDate: Boolean, hasTimestamp: Boolean): Boolean = { + + // Mirrors Spark's DataSourceUtils.datetimeRebaseSpec/int96RebaseSpec: when the file has no + // Spark version key the configured mode decides (EXCEPTION must also fall back, because + // Spark would raise on ancient values while Comet would not); when a version is present the + // mode is ignored and only the version and the legacy markers matter. The `v < minVersion` + // comparison is deliberately the same lexicographic string comparison Spark uses. + def needsRebase(facts: DatetimeFooterFacts): Boolean = { + def modeNeedsRebase(mode: String, minVersion: String, hasLegacyKey: Boolean): Boolean = + facts.sparkVersion.fold(mode != "CORRECTED")(v => v < minVersion || hasLegacyKey) + + (hasDate && + modeNeedsRebase(datetimeMode, "3.0.0", facts.hasLegacyDateTime)) || + (hasTimestamp && + (modeNeedsRebase(datetimeMode, "3.0.0", facts.hasLegacyDateTime) || + modeNeedsRebase(int96Mode, "3.1.0", facts.hasLegacyInt96))) + } + + def readFacts(path: Path): DatetimeFooterFacts = { + val inputFile = HadoopInputFile.fromPath(path, conf) + val readOptions = HadoopReadOptions + .builder(conf, path) + .withMetadataFilter(SKIP_ROW_GROUPS) + .build() + val reader = ParquetFileReader.open(inputFile, readOptions) + try { + val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData + DatetimeFooterFacts( + Option(metadata.get("org.apache.spark.version")), + metadata.containsKey("org.apache.spark.legacyDateTime"), + metadata.containsKey("org.apache.spark.legacyINT96")) + } finally { + reader.close() + } + } + + // Answer from the cache where possible; only cache misses pay a footer read. + var cachedNeedsRebase = false + val misses = new ListBuffer[(FooterCacheKey, Path)] + files.foreach { file => + val key = (file.path.toString, file.length, file.modificationTime) + val cached = footerFactsCache.get(key) + if (cached != null) { + cachedNeedsRebase = cachedNeedsRebase || needsRebase(cached) + } else { + misses += ((key, file.path)) + } + } + if (cachedNeedsRebase) { + return true + } + if (misses.isEmpty) { + return false + } + val parallelism = 8 val pool = ThreadUtils.newDaemonFixedThreadPool(parallelism, "checkingParquetDatetimeRebase") - val completion = new ExecutorCompletionService[Boolean](pool) - val remaining = paths.iterator + val completion = new ExecutorCompletionService[(FooterCacheKey, DatetimeFooterFacts)](pool) + val remaining = misses.iterator var inFlight = 0 // Spark's Parquet footer reader uses ThreadUtils.parmap, which submits every input eagerly. // Keep only `parallelism` reads in flight so finding one legacy footer stops further reads. def submitNext(): Unit = { - val path = remaining.next() - completion.submit(new Callable[Boolean] { - override def call(): Boolean = { - val inputFile = HadoopInputFile.fromPath(path, conf) - val readOptions = HadoopReadOptions - .builder(conf, path) - .withMetadataFilter(SKIP_ROW_GROUPS) - .build() - val reader = ParquetFileReader.open(inputFile, readOptions) - try { - val metadata = reader.getFooter.getFileMetaData.getKeyValueMetaData - val version = Option(metadata.get("org.apache.spark.version")) - def needsRebase(mode: String, minVersion: String, legacyKey: String): Boolean = - version.fold(mode != "CORRECTED")(v => - v < minVersion || metadata.containsKey(legacyKey)) - - (hasDate && - needsRebase(datetimeMode, "3.0.0", "org.apache.spark.legacyDateTime")) || - (hasTimestamp && - (needsRebase(datetimeMode, "3.0.0", "org.apache.spark.legacyDateTime") || - needsRebase(int96Mode, "3.1.0", "org.apache.spark.legacyINT96"))) - } finally { - reader.close() - } - } + val (key, path) = remaining.next() + completion.submit(new Callable[(FooterCacheKey, DatetimeFooterFacts)] { + override def call(): (FooterCacheKey, DatetimeFooterFacts) = (key, readFacts(path)) }) inFlight += 1 } @@ -85,7 +144,9 @@ object CometScanUtils { } var requiresRebase = false while (!requiresRebase && inFlight > 0) { - requiresRebase = completion.take().get() + val (key, facts) = completion.take().get() + footerFactsCache.put(key, facts) + requiresRebase = needsRebase(facts) inFlight -= 1 if (!requiresRebase && remaining.hasNext) { submitNext() diff --git a/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala index ae7766bfbac..5b1b710fdad 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala @@ -396,6 +396,43 @@ class CometParquetWriterSuite extends CometTestBase { } } + test("parquet write with LEGACY datetime rebase mode falls back to Spark") { + withTempPath { dir => + val df = spark.sql( + "SELECT id, date'1000-01-01' AS d, timestamp'1000-01-01 00:00:00' AS ts FROM range(10)") + + // The native writer always writes corrected (proleptic Gregorian) values, so a LEGACY + // write rebase mode must fall back to Spark, which rebases the values and stamps the + // legacy markers. + val legacyPath = new File(dir, "legacy.parquet").getAbsolutePath + withSQLConf( + CometConf.COMET_NATIVE_PARQUET_WRITE_ENABLED.key -> "true", + CometConf.getOperatorAllowIncompatConfigKey(classOf[DataWritingCommandExec]) -> "true", + CometConf.COMET_EXEC_ENABLED.key -> "true", + SQLConf.PARQUET_REBASE_MODE_IN_WRITE.key -> "LEGACY", + SQLConf.PARQUET_INT96_REBASE_MODE_IN_WRITE.key -> "LEGACY") { + + val plan = captureWritePlan(path => df.write.parquet(path), legacyPath) + assertNoCometNativeWriteExec(plan) + } + checkAnswer(spark.read.parquet(legacyPath), df.collect()) + + // The same write with corrected modes stays native. + val correctedPath = new File(dir, "corrected.parquet").getAbsolutePath + withSQLConf( + CometConf.COMET_NATIVE_PARQUET_WRITE_ENABLED.key -> "true", + CometConf.getOperatorAllowIncompatConfigKey(classOf[DataWritingCommandExec]) -> "true", + CometConf.COMET_EXEC_ENABLED.key -> "true", + SQLConf.PARQUET_REBASE_MODE_IN_WRITE.key -> "CORRECTED", + SQLConf.PARQUET_INT96_REBASE_MODE_IN_WRITE.key -> "CORRECTED") { + + val plan = captureWritePlan(path => df.write.parquet(path), correctedPath) + assertHasCometNativeWriteExec(plan) + } + checkAnswer(spark.read.parquet(correctedPath), df.collect()) + } + } + test("parquet write with temporal types within complex types") { withTempPath { dir => val outputPath = new File(dir, "output.parquet").getAbsolutePath diff --git a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala index 35bb56faeb7..86da18bbc6f 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -29,7 +29,7 @@ import scala.reflect.runtime.universe.TypeTag import org.scalactic.source.Position import org.scalatest.Tag -import org.apache.hadoop.fs.Path +import org.apache.hadoop.fs.{FileUtil, Path} import org.apache.parquet.example.data.simple.SimpleGroup import org.apache.parquet.schema.MessageTypeParser import org.apache.spark.SparkException @@ -1766,9 +1766,13 @@ class ParquetReadV1Suite extends ParquetReadSuite with AdaptiveSparkPlanHelper { .find(_.getName.endsWith(".parquet")) .get } - def requiresRebase(paths: Seq[Path]): Boolean = + def fileInfo(p: Path): CometScanUtils.ParquetFileInfo = { + val status = p.getFileSystem(hadoopConf).getFileStatus(p) + CometScanUtils.ParquetFileInfo(p, status.getLen, status.getModificationTime) + } + def requiresRebase(files: Seq[CometScanUtils.ParquetFileInfo]): Boolean = CometScanUtils.requiresDatetimeRebase( - paths, + files, hadoopConf, "CORRECTED", "CORRECTED", @@ -1777,10 +1781,19 @@ class ParquetReadV1Suite extends ParquetReadSuite with AdaptiveSparkPlanHelper { val correctedFile = parquetFile(correctedPath) val legacyFile = parquetFile(legacyPath) - assert(requiresRebase(Seq(correctedFile, legacyFile))) - assert( - requiresRebase( - Seq.fill(8)(legacyFile) :+ new Path(path.toString, "must-not-be-read.parquet"))) + assert(requiresRebase(Seq(fileInfo(correctedFile), fileInfo(legacyFile)))) + + // Early exit: eight not-yet-cached legacy footers fill every read slot, so the + // nonexistent ninth file must never be opened. + val fs = legacyFile.getFileSystem(hadoopConf) + val legacyCopies = (0 until 8).map { i => + val copy = new Path(path.toString, s"legacy-copy-$i.parquet") + FileUtil.copy(fs, legacyFile, fs, copy, false, hadoopConf) + fileInfo(copy) + } + val missingFile = + CometScanUtils.ParquetFileInfo(new Path(path.toString, "must-not-be-read.parquet"), 0, 0) + assert(requiresRebase(legacyCopies :+ missingFile)) val df = spark.read.parquet(correctedPath.toString, legacyPath.toString) val plan = df.queryExecution.executedPlan @@ -1794,6 +1807,21 @@ class ParquetReadV1Suite extends ParquetReadSuite with AdaptiveSparkPlanHelper { 1L, java.sql.Timestamp.valueOf("1000-01-01 00:00:00"), java.sql.Timestamp.valueOf("1000-01-01 00:00:00"))) + + // Disabling the check restores the previous behavior: the legacy file no longer causes + // a fallback (and its ancient values are then read without rebasing). + withSQLConf(CometConf.COMET_SCAN_PARQUET_CHECK_DATETIME_REBASE.key -> "false") { + val uncheckedPlan = + spark.read.parquet(legacyPath.toString).queryExecution.executedPlan + assert(collect(uncheckedPlan) { case _: CometNativeScanExec => true }.nonEmpty) + } + + // Footer facts are cached per (path, length, modificationTime): the corrected file can + // be answered again without any I/O even after the underlying file is deleted. + val correctedInfo = fileInfo(correctedFile) + assert(!requiresRebase(Seq(correctedInfo))) + fs.delete(correctedFile, false) + assert(!requiresRebase(Seq(correctedInfo))) } } @@ -1806,9 +1834,12 @@ class ParquetReadV1Suite extends ParquetReadSuite with AdaptiveSparkPlanHelper { |""".stripMargin).coalesce(1).write.parquet(path.toString) } + // NTZ values themselves are never rebased by Spark; the check only matters where a + // Parquet TIMESTAMP (LTZ/INT96) column may be read as NTZ, which Comet permits only + // when COMET_ALLOW_TIMESTAMP_LTZ_AS_NTZ is true (Spark 4.x). val (_, cometPlan) = checkSparkAnswer(spark.read.parquet(path.toString).select("ts_ntz")) val nativeScans = collect(cometPlan) { case _: CometNativeScanExec => true } - if (CometConf.COMET_SCHEMA_EVOLUTION_ENABLED) { + if (CometConf.COMET_ALLOW_TIMESTAMP_LTZ_AS_NTZ) { assert(nativeScans.isEmpty) } else { assert(nativeScans.nonEmpty)