andygrove commented on code in PR #5859: URL: https://github.com/apache/datafusion-comet/pull/5859#discussion_r4105059165
########## spark/src/test/scala/org/apache/spark/sql/benchmark/CometCacheRowReaderBenchmark.scala: ########## @@ -0,0 +1,207 @@ +/* + * 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.spark.sql.benchmark + +import java.nio.charset.StandardCharsets + +import org.apache.spark.benchmark.BenchmarkBase +import org.apache.spark.sql.{DataFrame, Row, SparkSession} +import org.apache.spark.sql.comet.execution.arrow.ArrowCachedBatchSerializer +import org.apache.spark.sql.execution.ColumnarToRowExec +import org.apache.spark.sql.execution.columnar.{CometInMemoryRelationHelper, DefaultCachedBatch, DefaultCachedBatchSerializer, InMemoryRelation, InMemoryTableScanExec} +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.storage.StorageLevel + +import org.apache.comet.{CometConf, CometSparkSessionExtensions} + +/** + * Compare Spark consumers of Comet and Spark caches (issue #5485). + * + * Arguments: [spark|comet|comet-row|all] [rows] [iterations] [all|mixed|numeric]. Run one format + * per JVM in alternating order. comet requires the fused columnar reader; comet-row disables + * vectorized cache reading to isolate the row iterator. Cache creation and validation are outside + * timing. + */ +object CometCacheRowReaderBenchmark extends BenchmarkBase { + private val warmups = 5 + + override def runBenchmarkSuite(args: Array[String]): Unit = { + require(args.length <= 4, "Expected format, rows, iterations, schema") + val format = args.headOption.getOrElse("all") + val rows = args.lift(1).map(_.toLong).getOrElse(5000000L) + val iterations = args.lift(2).map(_.toInt).getOrElse(15) + val schema = args.lift(3).getOrElse("all") + require(Set("all", "spark", "comet", "comet-row").contains(format)) + require(Set("all", "mixed", "numeric").contains(schema)) + require(rows > 0 && iterations > 0) + + emit("CACHE_SAMPLE,format,schema,query,rows,iteration,elapsed_ns") + val formats = + if (format == "all") Seq("spark", "comet", "comet-row") else Seq(format) + val schemas = if (schema == "all") Seq("mixed", "numeric") else Seq(schema) + formats.foreach { name => + CometInMemoryRelationHelper.clearSerializer() + SparkSession.clearActiveSession() + SparkSession.clearDefaultSession() + val serializer = if (name == "spark") { + classOf[DefaultCachedBatchSerializer].getName + } else { + classOf[ArrowCachedBatchSerializer].getName + } + val spark = SparkSession + .builder() + .master("local[1]") + .appName(getClass.getSimpleName) + .config("spark.ui.enabled", "false") + .config("spark.sql.cache.serializer", serializer) + .config("spark.sql.shuffle.partitions", "1") + .config("spark.sql.inMemoryColumnarStorage.batchSize", "10000") + .config("spark.sql.inMemoryColumnarStorage.compressed", "true") + .config("spark.io.compression.codec", "lz4") + .config(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key, "false") + .config(SQLConf.WHOLESTAGE_CODEGEN_ENABLED.key, "true") + .config(SQLConf.CACHE_VECTORIZED_READER_ENABLED.key, (name != "comet-row").toString) + .config(SQLConf.CODEGEN_FACTORY_MODE.key, "CODEGEN_ONLY") + .config(CometConf.COMET_ENABLED.key, "true") + .config(CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key, "true") + .config(CometConf.COMET_EXEC_ENABLED.key, "false") + .config(CometConf.COMET_SHUFFLE_ENABLED.key, "false") Review Comment: This session turns Comet on but never enables on-heap or off-heap memory for it, so `isCometLoaded` returns false and the rule never fires. Without `ENABLE_COMET_ONHEAP=true` in the environment, which `make benchmark-...` does not set, the `comet` arm fails its own `Expected the fused columnar cache reader` assertion. `CometBenchmarkBase` sets `spark.comet.exec.onHeap.enabled` for this reason. Could you set it here too? ########## spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala: ########## @@ -351,6 +352,161 @@ class CometInMemoryCacheSuite extends CometTestBase { } } + test("Spark row consumers of Comet cache preserve values across batches") { + for { + mode <- Seq("CODEGEN_ONLY", "NO_CODEGEN") + vectorized <- Seq(false, true) + } { + withSQLConf( + CometConf.COMET_ENABLED.key -> "false", Review Comment: This test still sets `spark.comet.enabled=false`, and since the rule now checks `isCometLoaded`, none of the four combinations take the fused path. The `vectorized && CODEGEN_ONLY` case runs the row iterator like the others, and the test only asserts that `ColumnarToRowExec` is absent, so nothing notices. That leaves `sum(key)` and `sum(length(s))` in the next test as the only fused-path coverage. Could this run with Comet on and exec and shuffle off, like the next test, and assert that the transition is present for `vectorized && CODEGEN_ONLY`? A consumer that reads every column through codegen, such as `df.filter($"key" >= 0)`, would cover the scalar types as well. I tried that locally and the whole type matrix passes through the fused path with AQE on and off. ########## spark/src/test/scala/org/apache/spark/sql/benchmark/CometCacheRowReaderBenchmark.scala: ########## @@ -0,0 +1,207 @@ +/* + * 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.spark.sql.benchmark + +import java.nio.charset.StandardCharsets + +import org.apache.spark.benchmark.BenchmarkBase +import org.apache.spark.sql.{DataFrame, Row, SparkSession} +import org.apache.spark.sql.comet.execution.arrow.ArrowCachedBatchSerializer +import org.apache.spark.sql.execution.ColumnarToRowExec +import org.apache.spark.sql.execution.columnar.{CometInMemoryRelationHelper, DefaultCachedBatch, DefaultCachedBatchSerializer, InMemoryRelation, InMemoryTableScanExec} +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.storage.StorageLevel + +import org.apache.comet.{CometConf, CometSparkSessionExtensions} + +/** + * Compare Spark consumers of Comet and Spark caches (issue #5485). + * + * Arguments: [spark|comet|comet-row|all] [rows] [iterations] [all|mixed|numeric]. Run one format + * per JVM in alternating order. comet requires the fused columnar reader; comet-row disables + * vectorized cache reading to isolate the row iterator. Cache creation and validation are outside + * timing. + */ +object CometCacheRowReaderBenchmark extends BenchmarkBase { Review Comment: #5543 added a Spark-operator arm to `CometInMemoryCacheBenchmark` after this PR was opened. Would it make sense to add the fused case there instead of a second harness? It already extends `CometBenchmarkBase`, and it is the benchmark the user guide's Limitations table cites. Its Spark-operator arm runs with Comet off, so after this PR it only measures the row iterator. The numbers in the description also predate #5543's storage format change, so could you re-run them on current main? ########## spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala: ########## @@ -351,6 +352,161 @@ class CometInMemoryCacheSuite extends CometTestBase { } } + test("Spark row consumers of Comet cache preserve values across batches") { + for { + mode <- Seq("CODEGEN_ONLY", "NO_CODEGEN") + vectorized <- Seq(false, true) + } { + withSQLConf( + CometConf.COMET_ENABLED.key -> "false", + SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", + SQLConf.CACHE_VECTORIZED_READER_ENABLED.key -> vectorized.toString, + SQLConf.COLUMN_BATCH_SIZE.key -> "7", + SQLConf.CODEGEN_FACTORY_MODE.key -> mode, + SQLConf.WHOLESTAGE_CODEGEN_ENABLED.key -> (mode == "CODEGEN_ONLY").toString, + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.SHUFFLE_PARTITIONS.key -> "2") { + val scalars = Seq( + "boolean", + "tinyint", + "smallint", + "int", + "bigint", + "float", + "double", + "decimal(10,2)", + "decimal(38,2)", + "date", + "timestamp", + "timestamp_ntz").zipWithIndex.map { case (dt, i) => + val value = dt match { + case "date" | "timestamp" | "timestamp_ntz" => + s"cast(date_add(DATE '2000-01-01', cast(id AS INT)) AS $dt)" + case _ => s"cast(id AS $dt)" + } + s"if(id % 3 = 0, null, $value) AS c$i" + } + val source = spark + .range(0, 41, 1, 2) + .selectExpr((Seq("id AS key") ++ scalars ++ Seq( + "if(id % 3 = 0, null, repeat(concat('字', id), cast(id + 1 AS INT))) AS s", + "if(id % 3 = 0, null, cast(concat('binary', id) AS BINARY)) AS b", + "if(id % 3 = 0, null, array(cast(id AS STRING), null)) AS a", + "if(id % 3 = 0, null, named_struct('x', id, 'a', array(cast(id AS STRING)))) AS st", + "if(id % 3 = 0, null, map('k', array(cast(id AS STRING), null))) AS m", + "null AS n")): _*) + + def queries(df: DataFrame): Seq[DataFrame] = Seq( + df.select("*"), + df.selectExpr("s AS renamed", "key", "b", "a", "st", "m"), + df.orderBy($"s".desc, $"key"), + df.join(spark.range(41).toDF("join_key"), $"key" === $"join_key").select(df("*")), + df.selectExpr("count(*)"), + df.limit(1)) + + val expected = queries(source).map(_.collect().toSeq) + source.cache() + try { + assert(source.count() == 41) + val relation = + spark.sharedState.cacheManager.lookupCachedData(source).get.cachedRepresentation + val buffers = relation.cacheBuilder.cachedColumnBuffers.collect() + assert(buffers.length > 2) + assert(buffers.forall(_.getClass.getSimpleName == "CometCachedBatch")) + queries(source).zip(expected).foreach { case (df, answer) => + val scans = + df.queryExecution.executedPlan.collect { case scan: InMemoryTableScanExec => + scan + } + assert( + scans.nonEmpty && scans.forall(_.supportsColumnar == vectorized), + df.queryExecution.executedPlan.toString) + if (!vectorized || mode == "NO_CODEGEN") { + assert( + !df.queryExecution.executedPlan.exists(_.isInstanceOf[ColumnarToRowExec]), + df.queryExecution.executedPlan.toString) + } + checkAnswer(df, answer) + } + } finally source.unpersist(blocking = true) + } + } + } + + test("Spark generated cache consumers respect runtime enable and codegen settings") { + for { + adaptive <- Seq(false, true) + disabledSetting <- Seq( + CometConf.COMET_ENABLED.key -> "false", + CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "false", + SQLConf.CODEGEN_FACTORY_MODE.key -> "NO_CODEGEN", + SQLConf.WHOLESTAGE_CODEGEN_ENABLED.key -> "false") + } { + withSQLConf( + CometConf.COMET_ENABLED.key -> "true", + CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "true", + CometConf.COMET_EXEC_ENABLED.key -> "false", + CometConf.COMET_SHUFFLE_ENABLED.key -> "false", + SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> adaptive.toString, + SQLConf.CACHE_VECTORIZED_READER_ENABLED.key -> "true", + SQLConf.WHOLESTAGE_CODEGEN_ENABLED.key -> "true", + SQLConf.CODEGEN_FACTORY_MODE.key -> "CODEGEN_ONLY", + SQLConf.COLUMN_BATCH_SIZE.key -> "7", + SQLConf.SHUFFLE_PARTITIONS.key -> "2") { + val source = spark + .range(0, 41, 1, 2) + .selectExpr("id AS key", "if(id % 3 = 0, null, concat('字', id)) AS s") + def query = source + .filter("key >= 7") + .selectExpr("sum(key)", "sum(length(s))", "count(*)") + val expected = query.collect().toSeq + source.cache() + try { + val builder = spark.sharedState.cacheManager + .lookupCachedData(source) + .get + .cachedRepresentation + .cacheBuilder + // Materialize with fusion enabled, then disable and re-enable it on the same cache. + Seq(true, false, true).zipWithIndex.foreach { case (enabled, index) => + val settings = if (enabled) Seq.empty else Seq(disabledSetting) + withSQLConf(settings: _*) { + val cold = index == 0 + val df = query + val plan = df.queryExecution.executedPlan + // Planning must not materialize the cache or replace AQE's cache-stage metadata. + assert(builder.isCachedColumnBuffersLoaded != cold, plan.toString) + checkAnswer(df, expected) Review Comment: `checkAnswer` first runs `df.materializedRdd.count()` on a separate query execution and only then `df.collect()`. So the cold iteration's cache is already loaded when the plan asserted here runs, and the path where the table-cache stage materializes and AQE re-plans is never the one checked. The #6208 test earlier in this file uses `QueryTest.checkAnswer(df, expected, checkToRDD = false)` for this reason. Could you do the same here? I tried it locally, and the cold run is then really cold and still gets one fused transition. ########## spark/src/main/scala/org/apache/comet/CometConf.scala: ########## @@ -263,20 +263,25 @@ object CometConf extends ShimCometConf { val COMET_EXEC_IN_MEMORY_CACHE_ENABLED: ConfigEntry[Boolean] = conf("spark.comet.exec.inMemoryCache.enabled") .category(CATEGORY_EXEC) - .doc("Whether to enable Comet native execution for in-memory cached tables. Its value at " + - "startup also decides whether CometDriverPlugin installs Comet's cache serializer, " + - "which stores cached data in Arrow format. Because spark.sql.cache.serializer is a " + - "static config, the cached format is fixed for the application, and disabling this " + - "at runtime only sends cached scans back to Spark's execution path. Relations whose " + - "schema Comet's Arrow writer does not support are always cached in Spark's default " + - "format. Each cached batch is stored as one Arrow IPC record batch with per-buffer " + - "zstd compression, and a scan copies out only the buffers of the columns it projected, " + - "so the unselected ones are never decompressed. Reads that feed Spark operators rather " + - "than Comet ones still pay a row conversion the default format avoids, and can be " + - "slower than Spark's cache. With spark.kryo.registrationRequired=true, also set " + - "spark.kryo.registrator=org.apache.comet.CometKryoRegistrator before creating the " + - "SparkContext, otherwise caching fails as soon as a block is serialized, including " + - "the disk half of the default MEMORY_AND_DISK storage level.") + .doc( + "Whether to enable Comet native scans and fused Spark reads of in-memory cached tables. " + + "Requires spark.comet.enabled=true. At startup, this setting also decides whether " + + "CometDriverPlugin installs Comet's cache serializer, which stores cached data in " + + "Arrow format. Because spark.sql.cache.serializer is a " + + "static config, the cached format is fixed for the application, and disabling this " + + "or spark.comet.enabled at runtime sends cached scans back to Spark's execution path " + + "without the fused reader. Relations whose schema Comet's Arrow writer does not " + + "support are always cached in Spark's default " + + "format. Each cached batch is stored as one Arrow IPC record batch with per-buffer " + + "zstd compression, and a scan copies out only the buffers of the columns it " + + "projected, so the unselected ones are never decompressed. Eligible Spark " + + "whole-stage codegen consumers read cached vectors directly when vectorized cache " + + "reading is enabled; other Spark row consumers use a reusable row buffer. Decoding " + + "costs can still make wide numeric reads slower than Spark's default cache. With " + Review Comment: This describes the new split, but the Limitations section of `docs/source/user-guide/latest/in-memory-cache.md` still says Spark-operator reads are slower than Spark's own format and that the cause is not yet established. Its table was measured with Comet off, and those reads now go through `CachedBatchRowIterator`. With Comet on and exec off, eligible consumers take the fused path instead. Could you update that section in this PR, including when the fused path applies? With exec on, a Spark operator above the cache already reads through `CometColumnarToRow` over the native scan, so that case does not change. ########## spark/src/main/scala/org/apache/comet/rules/CometCacheColumnarRule.scala: ########## @@ -0,0 +1,95 @@ +/* + * 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 org.apache.spark.sql.catalyst.expressions.LeafExpression +import org.apache.spark.sql.catalyst.expressions.codegen.CodegenFallback +import org.apache.spark.sql.catalyst.rules.Rule +import org.apache.spark.sql.comet.execution.arrow.ArrowCachedBatchSerializer +import org.apache.spark.sql.execution.{CodegenSupport, ColumnarToRowExec, ColumnarToRowTransition, SparkPlan, WholeStageCodegenExec} +import org.apache.spark.sql.execution.adaptive.QueryStageExec +import org.apache.spark.sql.execution.columnar.InMemoryTableScanExec +import org.apache.spark.sql.internal.SQLConf + +import org.apache.comet.CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED +import org.apache.comet.CometSparkSessionExtensions.isCometLoaded + +/** + * Lets Spark's generated consumers read cached Arrow vectors without an intermediate UnsafeRow. + * + * Data flows upward. Spark's InputAdapter/whole-stage wrappers and an optional AQE cache stage + * are omitted: + * {{{ + * Before After + * +------------------------+ +------------------------+ + * | Spark codegen consumer | | Spark codegen consumer | + * +------------------------+ +------------------------+ + * ^ ^ + * | UnsafeRow | column values + * +------------------------+ +------------------------+ + * | InMemoryTableScanExec | | ColumnarToRowExec | + * | row iterator | | fused with consumer | + * +------------------------+ +------------------------+ + * ^ + * | ColumnarBatch + * +------------------------+ + * | InMemoryTableScanExec | + * | Arrow vectors | + * +------------------------+ + * }}} + */ +object CometCacheColumnarRule extends Rule[SparkPlan] { + override def apply(plan: SparkPlan): SparkPlan = { + if (!isCometLoaded(conf) || !COMET_EXEC_IN_MEMORY_CACHE_ENABLED.get(conf)) return plan Review Comment: Since #5394, `CometColumnar.postColumnarTransitions` runs these rules in plan-only mode too, and this rule rewrites plans that contain no Comet operators. With `spark.comet.explain.planOnly.enabled=true`, a read of a Comet-format cache executes a `ColumnarToRowExec` that Spark would not have planned, while the plan-only doc says Spark executes the query unchanged. Should this also return early when plan-only is on? A plan-only case in the runtime-settings test would pin it down. -- 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]
