peterxcli opened a new issue, #5739:
URL: https://github.com/apache/datafusion-comet/issues/5739

   ## Problem
   
   Track the remaining filtered-scan correctness gap from #5553 and the 
nested-field statistics pruning needed before removing Comet's 
checked-timestamp safeguard.
   
   Comet currently sets `SparkParquetOptions.checked_timestamp_overflow = 
false` for **all scans with Spark data filters**, including filters Comet 
cannot serialize. An overflowing Parquet `TIMESTAMP_MILLIS` therefore becomes 
`NULL`, even when vanilla Spark reads the value and throws 
`ArithmeticException: long overflow`.
   
   The unfiltered nested-field fix can land independently. This issue covers 
the filtered-scan work that remains afterward.
   
   ## Upstream dependency
   
   https://github.com/apache/datafusion/issues/20871 tracks row-group 
statistics pruning for struct-field predicates. DataFusion 55 supports row 
filtering on primitive struct leaves, but does not yet use nested leaf 
statistics to prune row groups. These are distinct capabilities.
   
   Nested statistics pruning is a prerequisite, not a sufficient fix by itself: 
Comet must also preserve predicate semantics across millisecond-to-microsecond 
conversion and account for Spark's other pruning/filtering paths before 
enabling checked conversion for filtered scans.
   
   ## Verified reproduction
   
   Verified on Spark 4.1 with Comet's native DataFusion scan and DataFusion 
55.0.0, with Comet row-filter pushdown both off and on.
   
   Write one Parquet row group with plain encoding:
   
   | s.k | s.ts (raw TIMESTAMP_MILLIS integer) |
   | --- | --- |
   | 0 | 0 |
   | 1 | 9223372036854776 |
   
   ```sql
   SELECT s.k, s.ts FROM parquet_file WHERE s.k >= 0;
   ```
   
   Both rows satisfy the predicate, so statistics pruning cannot remove the 
overflowing row.
   
   - Vanilla Spark: throws `java.lang.ArithmeticException: long overflow`.
   - Comet: succeeds with `(0, 1970-01-01 00:00:00)` and `(1, NULL)`.
   
   The second timestamp fits in a signed 64-bit millisecond value, but 
multiplying it by 1,000 does not fit in a signed 64-bit microsecond value.
   
   ### Runnable Comet test
   
   Paste this test into 
`spark/src/test/scala/org/apache/comet/parquet/ParquetReadSuite.scala`. It uses 
the suite's existing Parquet writer and Spark/Comet comparison helpers. It 
intentionally asserts the current divergence, so a passing test reproduces the 
bug.
   
   ```scala
     test("reproduce #5553 filtered nested TIMESTAMP_MILLIS mismatch") {
       withTempDir { dir =>
         val path = new Path(dir.toURI.toString, "part-r-0.parquet")
         val schema = MessageTypeParser.parseMessageType("""
           |message root {
           |  optional group s {
           |    optional int32 k;
           |    optional int64 ts(TIMESTAMP_MILLIS);
           |  }
           |}
           |""".stripMargin)
         val writer = createParquetWriter(schema, path, dictionaryEnabled = 
false)
         Seq(0L, 9223372036854776L).zipWithIndex.foreach { case (millis, k) =>
           val record = new SimpleGroup(schema)
           val nested = record.addGroup(0)
           nested.add(0, k)
           nested.add(1, millis)
           writer.write(record)
         }
         writer.close()
   
         // Both rows pass the nested predicate, so Spark cannot prune the 
overflow away.
         // This deliberately asserts the remaining mismatch, rather than 
parity with Spark.
         Seq(false, true).foreach { rowFilterPushdown =>
           withSQLConf(
             SQLConf.PARQUET_VECTORIZED_READER_ENABLED.key -> "true",
             CometConf.COMET_PARQUET_ROW_FILTER_PUSHDOWN_ENABLED.key ->
               rowFilterPushdown.toString) {
             readParquetFile(path.toString) { df =>
               val query = df.where("s.k >= 0").select("s.k", "s.ts")
               assert(collect(query.queryExecution.executedPlan) { case _: 
CometNativeScanExec =>
                 true
               }.nonEmpty)
               val (sparkError, cometError) = checkSparkAnswerMaybeThrows(query)
               assert(sparkError.exists { error =>
                 Iterator.iterate(error)(_.getCause).takeWhile(_ != 
null).exists {
                   case _: ArithmeticException => true
                   case _ => false
                 }
               })
               assert(cometError.isEmpty)
               assert(query.collect().toSeq == Seq(Row(0, new 
java.sql.Timestamp(0)), Row(1, null)))
             }
           }
         }
       }
     }
   
   ```
   
   From the repository root, after building the native library:
   
   ```sh
   SPARK_LOCAL_IP=127.0.0.1 ./mvnw test -Dtest=none \
     -Dsuites="org.apache.comet.parquet.ParquetReadV1Suite reproduce #5553"
   ```
   
   ## Expected behavior and completion criteria
   
   - If Spark reads an overflowing value and throws, Comet must also throw, 
including for nested values.
   - If Spark prunes overflowing data and succeeds, Comet must also succeed.
   - Cover mixed/unprunable row groups and prunable row groups, nested 
timestamp predicates, and conversion-aware predicate rewriting.
   - Remove or narrow the filtered-scan safe-cast fallback only when those 
behaviors are verified; do not simply enable checked conversion globally.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to