diff --git a/docs/source/user-guide/latest/compatibility/scans.md b/docs/source/user-guide/latest/compatibility/scans.md index 57363125331..04247308812 100644 --- a/docs/source/user-guide/latest/compatibility/scans.md +++ b/docs/source/user-guide/latest/compatibility/scans.md @@ -43,6 +43,15 @@ 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. 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 @@ -51,17 +60,6 @@ The following features are not supported and cause Comet to fall back to Spark: - A read schema that repeats a Parquet field id, at the top level or within a struct, when `spark.sql.parquet.fieldId.read.enabled=true`. -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: - Selecting a field by name when multiple physical siblings match, including inside structs, diff --git a/native/core/src/execution/operators/parquet_writer.rs b/native/core/src/execution/operators/parquet_writer.rs index 9897ed626c6..a1bd2aa26f4 100644 --- a/native/core/src/execution/operators/parquet_writer.rs +++ b/native/core/src/execution/operators/parquet_writer.rs @@ -53,7 +53,7 @@ use futures::TryStreamExt; use parquet::{ arrow::ArrowWriter, basic::{Compression, GzipLevel, ZstdLevel}, - file::properties::WriterProperties, + file::{metadata::KeyValue, properties::WriterProperties}, }; use url::Url; @@ -240,6 +240,8 @@ pub struct ParquetWriterExec { column_names: Vec, /// Catalyst's target schema, including nullability and Parquet field metadata. output_schema: Option, + /// Runtime Spark version to record in the Parquet metadata + spark_version: String, /// Object store configuration options object_store_options: HashMap, /// Metrics @@ -261,6 +263,7 @@ impl ParquetWriterExec { partition_id: i32, column_names: Vec, output_schema: Option, + spark_version: String, object_store_options: HashMap, ) -> Result { // Preserve the input's partitioning so each partition writes its own file @@ -283,6 +286,7 @@ impl ParquetWriterExec { partition_id, column_names, output_schema, + spark_version, object_store_options, metrics: ExecutionPlanMetricsSet::new(), cache, @@ -471,6 +475,7 @@ impl ExecutionPlan for ParquetWriterExec { self.partition_id, self.column_names.clone(), self.output_schema.clone(), + self.spark_version.clone(), self.object_store_options.clone(), )?)), _ => Err(DataFusionError::Internal( @@ -531,6 +536,13 @@ 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.spark.version".to_string(), + Some(self.spark_version.clone()), + )])) .build(); let object_store_options = self.object_store_options.clone(); @@ -683,6 +695,7 @@ mod tests { 3, vec!["id".to_string()], None, + "4.2.0".to_string(), HashMap::new(), )?; @@ -753,6 +766,7 @@ mod tests { 0, vec!["required_id".to_string(), "values".to_string()], Some(output_schema), + "4.2.0".to_string(), HashMap::new(), )?; @@ -822,6 +836,7 @@ mod tests { 0, vec!["values".to_string()], Some(output_schema), + "4.2.0".to_string(), HashMap::new(), )?; @@ -1081,8 +1096,9 @@ mod tests { ParquetCompression::None, 0, // partition_id column_names, - None, // output_schema - HashMap::new(), // object_store_options + None, // output_schema + "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 10cc4ae37d7..e1544990818 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -2006,6 +2006,7 @@ impl PhysicalPlanner { writer.column_names.clone(), (!writer.output_schema.is_empty()) .then(|| convert_spark_types_to_arrow_schema(&writer.output_schema)), + writer.spark_version.clone(), object_store_options, )?); diff --git a/native/core/src/parquet/parquet_support.rs b/native/core/src/parquet/parquet_support.rs index 89221c51ebb..fa534a76a1b 100644 --- a/native/core/src/parquet/parquet_support.rs +++ b/native/core/src/parquet/parquet_support.rs @@ -91,8 +91,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 @@ -131,7 +129,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, @@ -148,7 +145,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/native/proto/src/proto/operator.proto b/native/proto/src/proto/operator.proto index 6a6284fb4c1..df3ecb177fc 100644 --- a/native/proto/src/proto/operator.proto +++ b/native/proto/src/proto/operator.proto @@ -928,6 +928,8 @@ message ParquetWriter { // Catalyst's target schema, including top-level nullability and Parquet field IDs. // Nested collection field IDs are carried by the corresponding DataType messages. repeated SparkStructField output_schema = 9; + // Runtime Spark version written to Parquet metadata for Spark reader compatibility. + string spark_version = 10; } enum AggregateMode { diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index c15e391c940..5237ae761df 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -991,6 +991,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 bb83297e636..304e697e17a 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 @@ -35,9 +36,10 @@ 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 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,6 +51,7 @@ import org.apache.comet.CometSparkSessionExtensions.{isCometLoaded, isSpark35Plu import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection} import org.apache.comet.objectstore.NativeConfig import org.apache.comet.parquet.CometParquetUtils.{encryptionEnabled, isEncryptionConfigSupported, readFieldId} +import org.apache.comet.serde.SupportLevel import org.apache.comet.serde.operator.{CometIcebergNativeScan, CometNativeScan} import org.apache.comet.shims.{CometTypeShim, ShimCometStreaming, ShimFileFormat, ShimSubqueryBroadcast} @@ -365,7 +368,45 @@ 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]) + // 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_ALLOW_TIMESTAMP_LTZ_AS_NTZ && + SupportLevel.containsType(scanExec.requiredSchema, classOf[TimestampNTZType])) + if ((hasDate || hasTimestamp) && COMET_SCAN_PARQUET_CHECK_DATETIME_REBASE.get()) { + val options = new ParquetOptions(r.options, conf) + val files = cometScan.selectedPartitions.iterator + .flatMap(_.files.iterator.map(f => + CometScanUtils.ParquetFileInfo(f.getPath, f.getLen, f.getModificationTime))) + .toSeq + try { + if (CometScanUtils.requiresDatetimeRebase( + files, + 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/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala index 6af2cc18ba6..e0a7faea22f 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 @@ -21,7 +21,7 @@ package org.apache.comet.serde.operator import scala.jdk.CollectionConverters._ -import org.apache.spark.SparkException +import org.apache.spark.{SPARK_VERSION_SHORT, SparkException} import org.apache.spark.sql.comet.{CometEmptyRelationExec, CometNativeExec, CometNativeWriteExec, CometScanWrapper} import org.apache.spark.sql.execution.SparkPlan import org.apache.spark.sql.execution.adaptive.QueryStageExec @@ -89,6 +89,10 @@ object CometDataWritingCommand extends CometOperatorSerde[DataWritingCommandExec return Unsupported(Some(s"Unsupported compression codec: $codec")) } + NativeWriteUtils + .legacyDatetimeRebaseWriteReason(cmd.query.output) + .foreach(reason => return Unsupported(Some(reason))) + Incompatible(Some("Parquet write support is highly experimental")) case _ => Unsupported(Some("Only Parquet writes are supported")) @@ -135,6 +139,7 @@ object CometDataWritingCommand extends CometOperatorSerde[DataWritingCommandExec .newBuilder() .setOutputPath(outputPath) .setCompression(codec) + .setSparkVersion(SPARK_VERSION_SHORT) .addAllColumnNames(cmd.query.output.map(_.name).asJava) .addAllOutputSchema(schema2Proto( cmd.query.schema.fields.toIndexedSeq, diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometWriteFiles.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometWriteFiles.scala index fffa5d7dc4e..fc6330a33d1 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometWriteFiles.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometWriteFiles.scala @@ -21,6 +21,7 @@ package org.apache.comet.serde.operator import org.apache.hadoop.conf.Configuration import org.apache.hadoop.fs.Path +import org.apache.spark.SPARK_VERSION_SHORT import org.apache.spark.sql.catalyst.util.CaseInsensitiveMap import org.apache.spark.sql.comet.{CometNativeExec, CometWriteFilesExec} import org.apache.spark.sql.execution.datasources.WriteFilesExec @@ -102,6 +103,10 @@ object CometWriteFiles extends CometOperatorSerde[WriteFilesExec] { return Unsupported(Some(s"Unsupported compression codec: $codec")) } + NativeWriteUtils + .legacyDatetimeRebaseWriteReason(op.child.output) + .foreach(reason => return Unsupported(Some(reason))) + Incompatible(Some("Parquet write support is highly experimental")) } @@ -135,6 +140,7 @@ object CometWriteFiles extends CometOperatorSerde[WriteFilesExec] { val writerOpBuilder = OperatorOuterClass.ParquetWriter .newBuilder() .setCompression(codec) + .setSparkVersion(SPARK_VERSION_SHORT) // getSupportLevel already declined the write if the tag is absent, so this cannot be empty. outputPathOf(op).foreach { outputPath => diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/NativeWriteUtils.scala b/spark/src/main/scala/org/apache/comet/serde/operator/NativeWriteUtils.scala index d8247e99a7f..3233fbfdbd9 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/NativeWriteUtils.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/NativeWriteUtils.scala @@ -23,11 +23,13 @@ import java.util.Locale import org.apache.hadoop.fs.Path import org.apache.parquet.hadoop.ParquetOutputFormat +import org.apache.spark.sql.catalyst.expressions.Attribute import org.apache.spark.sql.catalyst.plans.QueryPlan import org.apache.spark.sql.catalyst.util.CaseInsensitiveMap import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.types.{DateType, TimestampType} -import org.apache.comet.serde.OperatorOuterClass +import org.apache.comet.serde.{OperatorOuterClass, SupportLevel} import org.apache.comet.serde.QueryPlanSerde.serializeDataType /** @@ -186,6 +188,37 @@ object NativeWriteUtils { "spark.comet.parquet.write.enabled=false to write this table with Spark.") } + /** + * A fallback reason when the session asks for a LEGACY datetime rebase on write, or `None`. + * + * 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. + */ + def legacyDatetimeRebaseWriteReason(output: Seq[Attribute]): Option[String] = { + val hasDate = output.exists(a => SupportLevel.containsType(a.dataType, classOf[DateType])) + val hasTimestamp = + 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.isEmpty) { + None + } else { + Some( + "Native Parquet write always writes corrected (proleptic Gregorian) datetime values " + + s"and does not support LEGACY rebase mode (${legacyModeKeys.mkString(", ")})") + } + } + /** Compression codecs Comet's native Parquet writer can produce. */ val supportedCompressionCodecs: Set[String] = Set("none", "uncompressed", "snappy", "lz4", "zstd", "gzip") 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..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 @@ -19,11 +19,145 @@ 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 +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 { + /** 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( + 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[(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 (key, path) = remaining.next() + completion.submit(new Callable[(FooterCacheKey, DatetimeFooterFacts)] { + override def call(): (FooterCacheKey, DatetimeFooterFacts) = (key, readFacts(path)) + }) + inFlight += 1 + } + + try { + while (inFlight < parallelism && remaining.hasNext) { + submitNext() + } + var requiresRebase = false + while (!requiresRebase && inFlight > 0) { + val (key, facts) = completion.take().get() + footerFactsCache.put(key, facts) + requiresRebase = needsRebase(facts) + 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/CometParquetWriterSuite.scala b/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala index c46524e2917..b9493e0cb82 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala @@ -32,6 +32,7 @@ import org.apache.parquet.hadoop.ParquetFileReader import org.apache.parquet.hadoop.metadata.CompressionCodecName import org.apache.parquet.hadoop.util.HadoopInputFile import org.apache.parquet.schema.{MessageType, Type} +import org.apache.spark.SPARK_VERSION_SHORT import org.apache.spark.internal.io.FileCommitProtocol import org.apache.spark.sql.{AnalysisException, DataFrame, Row, SaveMode} import org.apache.spark.sql.catalyst.InternalRow @@ -672,6 +673,39 @@ class CometParquetWriterSuite extends CometParquetWriterTestBase { } } + 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 + withNativeWriter { + withSQLConf( + 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 + withNativeWriter { + withSQLConf( + 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 @@ -1540,6 +1574,16 @@ class CometParquetWriterSuite extends CometParquetWriterTestBase { // 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) 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 287d4ecb14a..066f82579bb 100644 --- a/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala @@ -31,7 +31,7 @@ import org.scalactic.source.Position import org.scalatest.Tag import org.apache.arrow.vector.types.pojo.{ArrowType, DictionaryEncoding, Field => ArrowField, FieldType, Schema => ArrowSchema} -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.hadoop.example.ExampleParquetWriter import org.apache.parquet.io.api.Binary @@ -40,7 +40,7 @@ import org.apache.spark.SparkException import org.apache.spark.sql.{CometTestBase, DataFrame, Row} import org.apache.spark.sql.catalyst.optimizer.{ConvertToLocalRelation, OptimizeIn} 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 @@ -2444,25 +2444,139 @@ 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). - 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) - - // Verify Comet scan is in the plan - val plan = df.queryExecution.executedPlan - checkCometOperators(plan) + 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) + } + } + } + } - // 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 => - val date = row.getDate(0) - assert(date.toLocalDate.getYear < 1582, s"Expected date before 1582, got $date") + 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).coalesce(1).write.parquet(legacyPath.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 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( + files, + hadoopConf, + "CORRECTED", + "CORRECTED", + hasDate = true, + hasTimestamp = true) + + val correctedFile = parquetFile(correctedPath) + val legacyFile = parquetFile(legacyPath) + 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 + 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( + 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))) + } + } + + 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) + } + + // 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_ALLOW_TIMESTAMP_LTZ_AS_NTZ) { + assert(nativeScans.isEmpty) + } else { + assert(nativeScans.nonEmpty) + } } } 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 cdd466f03d3..ec4de177bcf 100644 --- a/spark/src/test/scala/org/apache/spark/sql/CometTestBase.scala +++ b/spark/src/test/scala/org/apache/spark/sql/CometTestBase.scala @@ -825,6 +825,8 @@ abstract class CometTestBase .withPageSize(pageSize) .withDictionaryPageSize(dictionaryPageSize) .withPageRowCountLimit(pageRowCountLimit) + .withExtraMetaData( + java.util.Collections.singletonMap("org.apache.spark.version", SPARK_VERSION)) .withConf(hadoopConf) .build() }