comphead commented on code in PR #6753: URL: https://github.com/apache/datafusion-comet/pull/6753#discussion_r4231745645
########## spark/src/test/scala/org/apache/comet/rules/CometMultiStoreScanSuite.scala: ########## @@ -0,0 +1,382 @@ +/* + * 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.catalyst.plans.physical.UnknownPartitioning +import org.apache.spark.sql.comet.{CometCsvNativeScanExec, CometNativeScanExec, CometScanExec} +import org.apache.spark.sql.execution.{ExtendedMode, FileSourceScanExec, FormattedMode, 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.execution.exchange.ShuffleExchangeLike +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.types.{IntegerType, StructType} + +import org.apache.comet.CometConf +import org.apache.comet.CometConf.COMET_S3_COMPLIANT_SCHEMES_KEY +import org.apache.comet.hadoop.fs.FakeHdfsAuthorityFileSystem +import org.apache.comet.objectstore.NativeConfig + +/** + * 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) + + // Over the files `writeMixedLayout` writes (120, 15 and 10 bytes, with a 140 byte split), Spark + // packs the first two into one partition and the third into a second one. + private val twoPartitionsOneMixed = Seq( + SQLConf.FILES_MIN_PARTITION_NUM.key -> "1", + SQLConf.FILES_OPEN_COST_IN_BYTES.key -> "1", + SQLConf.FILES_MAX_PARTITION_BYTES.key -> "140") + + private val idSchema = new StructType().add("id", IntegerType) + + 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 storeKey(path: String): String = + NativeConfig.objectStoreKey(new URI(path), Set.empty, Set("hdfs")).key + + // A forwarded object store option that the plan's text must not show. + private val forwardedMarker = COMET_S3_COMPLIANT_SCHEMES_KEY -> "zzforwardedmarker" + + private def assertPlanHidesForwardedOptions(df: DataFrame): Unit = { + val texts = Seq( + df.queryExecution.executedPlan.toString, + df.queryExecution.explainString(ExtendedMode), + df.queryExecution.explainString(FormattedMode)) + texts.foreach(text => assert(!text.contains(forwardedMarker._2), text)) + } + + private def withoutComet(f: => Unit): Unit = + withSQLConf(CometConf.COMET_ENABLED.key -> "false")(f) + + private def writeIds( + path: String, + from: Int, + format: String = "parquet", + count: Int = 5): Unit = + withoutComet { + spark + .range(from.toLong, from.toLong + count) + .selectExpr("cast(id as int) as id") + .coalesce(1) + .write + .mode(SaveMode.Overwrite) + .format(format) + .save(path) + } + + /** CSV files of 120, 15 and 10 bytes at `first`, `second` and `third`. */ + private def writeMixedLayout(first: String, second: String, third: String): Unit = { + writeIds(first, 100, "csv", count = 30) + writeIds(second, 10, "csv") + writeIds(third, 0, "csv") + } + + /** 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) Review Comment: At this head a non-bucketed `CometNativeScanExec` defines `outputPartitioning` as `UnknownPartitioning(perPartitionData.length)` (`CometNativeScanExec.scala:141`), so this assert and the one at line 378 hold by construction, whatever the packing does. #6268 (open, approved) changes that to `UnknownPartitioning(0)` and sizes execution from `perPartitionData`, and it merges cleanly with this branch. From reading the code I expect both asserts to fail once the two meet, whichever lands second. I have not run it. Could we drop these two Parquet asserts? The CSV ones at lines 219 and 278 still check the new `CometCsvNativeScanExec.outputPartitioning`. -- 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]
