andygrove commented on code in PR #6753:
URL: https://github.com/apache/datafusion-comet/pull/6753#discussion_r4211101670


##########
spark/src/main/scala/org/apache/comet/serde/operator/CometNativeScan.scala:
##########
@@ -240,6 +240,22 @@ object CometNativeScan extends 
CometOperatorSerde[CometScanExec] with CometTypeS
       scan: CometScanExec,
       builder: Operator.Builder,
       childOp: OperatorOuterClass.Operator*): 
Option[OperatorOuterClass.Operator] = {
+    val hadoopConf =
+      
scan.relation.sparkSession.sessionState.newHadoopConfWithOptions(scan.relation.options)
+    // The root paths can miss a file, e.g. a catalog partition located 
outside the table, so
+    // check the listed files. The static partitions are a superset of what 
DPP keeps.
+    val multiStoreReason = CometScanUtils.multiStoreFallbackReason(

Review Comment:
   The scheme-family fallback turns off native execution for scans that 1.1 ran 
natively, such as a Hive table with some partitions on HDFS and some on S3. Now 
that files are packed per store, the only reason left is that `convert` 
forwards `extractObjectStoreOptions` for the first file's scheme. The keys for 
different families don't overlap, and native only reads the S3 and Azure ones 
(libhdfs and the `parse_url` path for `gs` ignore them). Could we forward the 
options of every scheme the scan reads, and keep the fallback only for an alias 
next to `s3` or `s3a` in the same bucket, where the settings really do collide?



##########
spark/src/test/scala/org/apache/comet/rules/CometMultiStoreScanSuite.scala:
##########
@@ -0,0 +1,290 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.comet.rules
+
+import java.io.File
+import java.net.URI
+import java.nio.file.Files
+import java.util.UUID
+
+import scala.jdk.CollectionConverters._
+
+import org.apache.commons.io.FileUtils
+import org.apache.spark.SparkConf
+import org.apache.spark.sql.{CometTestBase, DataFrame, SaveMode}
+import org.apache.spark.sql.catalyst.expressions.DynamicPruningExpression
+import org.apache.spark.sql.comet.{CometCsvNativeScanExec, 
CometNativeScanExec, CometScanExec}
+import org.apache.spark.sql.execution.{FileSourceScanExec, SparkPlan}
+import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper
+import org.apache.spark.sql.execution.datasources.FilePartition
+import org.apache.spark.sql.execution.datasources.v2.BatchScanExec
+import org.apache.spark.sql.internal.SQLConf
+import org.apache.spark.sql.types.{IntegerType, StructType}
+
+import org.apache.comet.CometConf
+import org.apache.comet.hadoop.fs.FakeHdfsAuthorityFileSystem
+
+/**
+ * Native scans over files in more than one object store, without a cloud 
store: `hdfs://nn1` and
+ * `hdfs://nn2` are two native stores backed by the local disk. Native 
execution cannot read them,
+ * so claimed scans are checked on their plans, and declined scans are run by 
Spark.
+ */
+class CometMultiStoreScanSuite extends CometTestBase with 
AdaptiveSparkPlanHelper {
+
+  private var rootDir: File = _
+
+  override protected def sparkConf: SparkConf = {
+    val conf = super.sparkConf
+    conf.set("spark.hadoop.fs.hdfs.impl", 
classOf[FakeHdfsAuthorityFileSystem].getName)
+    conf.set("spark.hadoop.fs.hdfs.impl.disable.cache", "true")
+    conf
+  }
+
+  override def beforeAll(): Unit = {
+    rootDir = 
Files.createTempDirectory(s"comet_multi_store_${UUID.randomUUID()}").toFile
+    super.beforeAll()
+  }
+
+  protected override def afterAll(): Unit = {
+    if (rootDir != null) FileUtils.deleteDirectory(rootDir)
+    super.afterAll()
+  }
+
+  private val nativeScan = Seq(
+    CometConf.COMET_NATIVE_SCAN_ENABLED.key -> "true",
+    CometConf.COMET_EXEC_ENABLED.key -> "true")
+
+  private val nativeCsv =
+    Seq(CometConf.COMET_CSV_V2_NATIVE_ENABLED.key -> "true", 
SQLConf.USE_V1_SOURCE_LIST.key -> "")
+
+  // Spark packs every file of the scan into one partition.
+  private val onePartition = Seq(
+    SQLConf.FILES_MIN_PARTITION_NUM.key -> "1",
+    SQLConf.FILES_OPEN_COST_IN_BYTES.key -> "1",
+    SQLConf.FILES_MAX_PARTITION_BYTES.key -> (128L * 1024 * 1024).toString)
+
+  private def hdfs(nameNode: String, name: String): String =
+    s"hdfs://$nameNode${rootDir.getAbsolutePath}/$name"
+
+  private def local(name: String): String = 
s"file://${rootDir.getAbsolutePath}/$name"
+
+  private def storeOf(path: String): String = new URI(path).getAuthority
+
+  private def withoutComet(f: => Unit): Unit =
+    withSQLConf(CometConf.COMET_ENABLED.key -> "false")(f)
+
+  private def writeIds(path: String, from: Int, format: String = "parquet"): 
Unit =
+    withoutComet {
+      spark
+        .range(from.toLong, from.toLong + 5)
+        .selectExpr("cast(id as int) as id")
+        .coalesce(1)
+        .write
+        .mode(SaveMode.Overwrite)
+        .format(format)
+        .save(path)
+    }
+
+  /** The files of each partition of Spark's own scans in `df`, which is 
planned, not run. */
+  private def sparkLayout(df: => DataFrame): Seq[Seq[String]] = {
+    var layout: Seq[Seq[String]] = Nil
+    withoutComet {
+      val plan = df.queryExecution.executedPlan
+      val partitions = collect(plan) {
+        case scan: FileSourceScanExec => scan.inputRDD.partitions.toSeq
+        case scan: BatchScanExec => scan.inputPartitions
+      }.flatten
+      layout = partitions.collect { case p: FilePartition =>
+        p.files.map(_.filePath.toString).toSeq
+      }
+    }
+    layout
+  }
+
+  /** The scans CometScanRule claims in the Spark plan of `df`. */
+  private def ruleClaims(df: => DataFrame): Seq[CometScanExec] = {
+    var sparkPlan: SparkPlan = null
+    withoutComet {
+      sparkPlan = df.queryExecution.executedPlan
+    }
+    CometScanRule(spark).apply(stripAQEPlan(sparkPlan)).collect { case s: 
CometScanExec => s }
+  }
+
+  private def nativeParquetScan(df: DataFrame): CometNativeScanExec = {
+    val plan = df.queryExecution.executedPlan
+    val scans = collect(plan) { case scan: CometNativeScanExec => scan }
+    assert(scans.size == 1, s"expected one native Parquet scan:\n$plan")
+    scans.head
+  }
+
+  test("parquet scan over two name nodes packs each name node's files on its 
own") {
+    val (a, b) = (hdfs("nn1", "two-nn-a"), hdfs("nn2", "two-nn-b"))
+    writeIds(a, 0)
+    writeIds(b, 10)
+    withSQLConf(nativeScan ++ onePartition: _*) {
+      val sparkFiles = sparkLayout(spark.read.parquet(a, b))
+      assert(sparkFiles.exists(_.map(storeOf).distinct.size > 1), s"Spark's: 
$sparkFiles")
+      val scan = nativeParquetScan(spark.read.parquet(a, b))
+      val cometFiles = scan.perPartitionFilePaths.toSeq
+      assert(cometFiles.forall(_.map(storeOf).distinct.size == 1), s"Comet's: 
$cometFiles")
+      assert(cometFiles.flatten.sorted == sparkFiles.flatten.sorted)
+      assert(cometFiles.size != sparkFiles.size)
+      assert(scan.outputPartitioning.numPartitions == 
scan.perPartitionData.length)
+    }
+  }
+
+  test("parquet scan over one name node keeps Spark's layout") {
+    val (a, b) = (hdfs("nn1", "one-nn-a"), hdfs("nn1", "one-nn-b"))
+    writeIds(a, 0)
+    writeIds(b, 10)
+    withSQLConf(nativeScan ++ onePartition: _*) {
+      val sparkFiles = sparkLayout(spark.read.parquet(a, b))
+      val scan = nativeParquetScan(spark.read.parquet(a, b))
+      assert(scan.perPartitionFilePaths.toSeq == sparkFiles)
+    }
+  }
+
+  test("csv scan over two name nodes splits partitions that mix name nodes") {

Review Comment:
   This test passes on 3.4 even though the plan above the scan is wrong there, 
because it only looks at the scan. Could we add a case with a global aggregate 
over a scan that Spark packs into one partition, and check that Comet either 
falls back or keeps the partition count? It can stay plan-only like the rest of 
the suite.



##########
spark/src/main/scala/org/apache/spark/sql/comet/CometCsvNativeScanExec.scala:
##########
@@ -92,7 +94,27 @@ object CometCsvNativeScanExec extends 
CometOperatorSerde[CometBatchScanExec] {
       val timeZone = sessionState.conf.sessionLocalTimeZone
       new CSVOptions(csvScan.options.asScala.toMap, columnPruning, timeZone)
     }
-    val filePartitions = op.inputPartitions.map(_.asInstanceOf[FilePartition])
+    val hadoopConf =
+      
sessionState.newHadoopConfWithOptions(op.session.sparkContext.conf.getAll.toMap)
+    val s3CompliantSchemes = NativeConfig.resolveS3CompliantSchemes(hadoopConf)
+    val libhdfsSchemes = NativeConfig.resolveLibhdfsSchemes(hadoopConf)
+    val inputPartitions = op.inputPartitions.map(_.asInstanceOf[FilePartition])
+    val multiStoreReason = CometScanUtils.multiStoreFallbackReason(
+      "Native CSV scan",
+      inputPartitions.view.flatMap(_.files.view.map(_.pathUri)),
+      s3CompliantSchemes,
+      libhdfsSchemes,
+      isBucketedScan = false)
+    if (multiStoreReason.nonEmpty) {
+      // CometExecRule falls back to the wrapped scan, so tag it too for the 
explain output.
+      withFallbackReason(op, multiStoreReason.get)
+      withFallbackReason(op.wrapped, multiStoreReason.get)
+      return None
+    }
+    // Native planning reads each partition through one object store.
+    val filePartitions = CometScanUtils.splitPartitionsByStore(

Review Comment:
   On Spark 3.4 this split can return wrong results. In 3.4, 
`DataSourceV2ScanExecBase.outputPartitioning` is `SinglePartition` when the 
scan has one input partition, so `EnsureRequirements` plans no exchange between 
the partial and final aggregate, or between the two sides of a sort merge join. 
When that one partition mixes stores, the split makes it two and everything 
above the scan runs on two partitions. I ran it on 3.4.3 with native CSV over 
two name nodes, using your `FakeHdfsAuthorityFileSystem` with the FS cache left 
on so libhdfs can read them. `agg(count(lit(1)), sum("id"))` returned `[5,10]` 
and `[5,60]` where Spark returns `[10,70]`. `count()` returned 5 instead of 10, 
and a join with the split scan on the right returned 5 of the 10 rows. With the 
split scan on the left, the join fails with `CsvScan has no file partition 1 (1 
partitions)`. 3.5 and later don't report `SinglePartition` there, so 4.1 is 
fine. Could the conversion fall back, or keep one partition, when the sp
 lit changes the partition count and `op.wrapped.outputPartitioning` is not an 
`UnknownPartitioning`? That check wouldn't need a shim.



-- 
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