From ed8ba0c95b6e7fc628d77ecf88118f15fb2cd947 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Sat, 10 Oct 2026 09:22:43 -0600 Subject: [PATCH 1/2] perf: generate the columnar-to-row projection once per executor CometColumnarToRowExec.doExecute built its UnsafeProjection with UnsafeProjection.create in every partition, which regenerates the projection's Java source each time. The plan runs that path whenever it is outside whole-stage codegen, which a schema with more nested fields than spark.sql.codegen.maxFields always is. For 100 columns nested three levels deep, generating the source took about 19 ms per partition, half the time of a 201-partition JVM shuffle of 1,000 rows. Cache the compiled class per executor by column layout and create a new instance for each partition. --- .../sql/comet/CometColumnarToRowExec.scala | 3 +- .../sql/comet/CometUnsafeProjection.scala | 146 +++++++++++++++ .../comet/CometUnsafeProjectionSuite.scala | 177 ++++++++++++++++++ 3 files changed, 325 insertions(+), 1 deletion(-) create mode 100644 spark/src/main/scala/org/apache/spark/sql/comet/CometUnsafeProjection.scala create mode 100644 spark/src/test/scala/org/apache/spark/sql/comet/CometUnsafeProjectionSuite.scala diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometColumnarToRowExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometColumnarToRowExec.scala index da1d9ba296..089ad6fbf6 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometColumnarToRowExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometColumnarToRowExec.scala @@ -77,7 +77,8 @@ case class CometColumnarToRowExec(child: SparkPlan) // plan (this) in the closure. val localOutput = this.output child.executeColumnar().mapPartitionsInternal { batches => - val toUnsafe = UnsafeProjection.create(localOutput, localOutput) + // Outside whole-stage codegen this runs once per partition, so reuse the generated class. + val toUnsafe = CometUnsafeProjection.create(localOutput) batches.flatMap { batch => numInputBatches += 1 numOutputRows += batch.numRows().toLong diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometUnsafeProjection.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometUnsafeProjection.scala new file mode 100644 index 0000000000..d787460a8b --- /dev/null +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometUnsafeProjection.scala @@ -0,0 +1,146 @@ +/* + * 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.comet + +import java.util.{LinkedHashMap => JLinkedHashMap, Map => JMap} +import java.util.concurrent.atomic.AtomicLong + +import org.apache.spark.sql.catalyst.expressions.{Attribute, BoundReference, CodeGeneratorWithInterpretedFallback, InterpretedUnsafeProjection, UnsafeProjection} +import org.apache.spark.sql.catalyst.expressions.codegen.{CodeAndComment, CodeFormatter, CodegenContext, CodeGenerator, GeneratedClass, GenerateUnsafeProjection} +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.types.DataType + +/** + * Creates the projection that copies each row of a batch into an UnsafeRow, as + * `UnsafeProjection.create(output, output)` does, but generates and compiles its class only once + * per executor for each column layout. + * + * `UnsafeProjection.create` generates the projection's Java source on every call, and only the + * compiled class is cached. `CometColumnarToRowExec` calls it once per partition whenever it runs + * outside whole-stage codegen, which Spark skips for a schema with more fields than + * `spark.sql.codegen.maxFields`, nested fields included. Spark's own scans fall back to rows for + * such schemas, so its `ColumnarToRowExec` rarely meets one, but Comet's operators always produce + * batches. The source grows with the number of nested fields, and each level of nesting repeats + * the work of splitting its children into methods: for 100 columns nested three levels deep, + * generating it took about 19 ms in every partition. + * + * Every call returns a new instance of the shared class, so callers own their projection as they + * do one from `UnsafeProjection.create`. + */ +private[comet] object CometUnsafeProjection + extends CodeGeneratorWithInterpretedFallback[Seq[BoundReference], UnsafeProjection] { + + /** + * What the generated source depends on: the type and nullability of each column, and the method + * size at which `CodegenContext.splitExpressions` splits the field writes. + */ + private case class Layout(columns: Seq[(DataType, Boolean)], methodSplitThreshold: Int) + + /** The same bound as Spark's compiled class cache, `spark.sql.codegen.cache.maxEntries`. */ + private[comet] val MaxCachedClasses = 100 + + /** Least recently used first. Guarded by `classes.synchronized`. */ + private val classes = new JLinkedHashMap[Layout, GeneratedClass](16, 0.75f, true) { + override def removeEldestEntry(eldest: JMap.Entry[Layout, GeneratedClass]): Boolean = + size() > MaxCachedClasses + } + + private val classesGenerated = new AtomicLong(0) + + /** How many projection classes this executor has generated, for tests. */ + private[comet] def generatedClassCount: Long = classesGenerated.get() + + /** A projection of rows with the columns of `output` to UnsafeRows. */ + def create(output: Seq[Attribute]): UnsafeProjection = + createObject(output.zipWithIndex.map { case (attr, ordinal) => + BoundReference(ordinal, attr.dataType, attr.nullable) + }) + + override protected def createCodeGeneratedObject( + columns: Seq[BoundReference]): UnsafeProjection = { + val layout = + Layout(columns.map(c => (c.dataType, c.nullable)), SQLConf.get.methodSplitThreshold) + val cached = classes.synchronized(classes.get(layout)) + if (cached != null) { + cached.generate(Array.empty[Any]).asInstanceOf[UnsafeProjection] + } else { + // Generated outside the lock: tasks that miss on the same layout at once each generate + // it, as every task does with UnsafeProjection.create. + val (generated, references) = generate(columns) + // A stored class is instantiated without references, so store only one that needs none. + // GenerateUnsafeProjection references no objects for bound columns. + if (references.isEmpty) { + val _ = classes.synchronized(classes.putIfAbsent(layout, generated)) + } + generated.generate(references).asInstanceOf[UnsafeProjection] + } + } + + override protected def createInterpretedObject(columns: Seq[BoundReference]): UnsafeProjection = + InterpretedUnsafeProjection.createProjection(columns) + + /** + * Generates and compiles the class that `GenerateUnsafeProjection.create` does, from the same + * template, which is repeated here because that method returns only an instance. Returns the + * class with the objects its instances reference. + */ + private def generate(columns: Seq[BoundReference]): (GeneratedClass, Array[Any]) = { + val ctx = new CodegenContext + val eval = GenerateUnsafeProjection.createCode(ctx, columns) + val body = + s""" + |public java.lang.Object generate(Object[] references) { + | return new SpecificUnsafeProjection(references); + |} + | + |class SpecificUnsafeProjection extends ${classOf[UnsafeProjection].getName} { + | + | private Object[] references; + | ${ctx.declareMutableStates()} + | + | public SpecificUnsafeProjection(Object[] references) { + | this.references = references; + | ${ctx.initMutableStates()} + | } + | + | public void initialize(int partitionIndex) { + | ${ctx.initPartition()} + | } + | + | // Scala.Function1 need this + | public java.lang.Object apply(java.lang.Object row) { + | return apply((InternalRow) row); + | } + | + | public UnsafeRow apply(InternalRow ${ctx.INPUT_ROW}) { + | ${eval.code} + | return ${eval.value}; + | } + | + | ${ctx.declareAddedFunctions()} + |} + """.stripMargin + val code = CodeFormatter.stripOverlappingComments( + new CodeAndComment(body, ctx.getPlaceHolderToComments())) + val (generatedClass, _) = CodeGenerator.compile(code) + classesGenerated.incrementAndGet() + (generatedClass, ctx.references.toArray) + } +} diff --git a/spark/src/test/scala/org/apache/spark/sql/comet/CometUnsafeProjectionSuite.scala b/spark/src/test/scala/org/apache/spark/sql/comet/CometUnsafeProjectionSuite.scala new file mode 100644 index 0000000000..07c6371f3e --- /dev/null +++ b/spark/src/test/scala/org/apache/spark/sql/comet/CometUnsafeProjectionSuite.scala @@ -0,0 +1,177 @@ +/* + * 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.comet + +import java.sql.Timestamp +import java.util.UUID + +import org.apache.spark.sql.{CometTestBase, Row} +import org.apache.spark.sql.catalyst.{CatalystTypeConverters, InternalRow} +import org.apache.spark.sql.catalyst.expressions.{AttributeReference, InterpretedUnsafeProjection, UnsafeProjection} +import org.apache.spark.sql.execution.WholeStageCodegenExec +import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.types._ + +class CometUnsafeProjectionSuite extends CometTestBase with AdaptiveSparkPlanHelper { + + /** A schema with every writer GenerateUnsafeProjection nests: structs, arrays and maps. */ + private def schema(nestedName: String = "a"): StructType = StructType( + Seq( + StructField("i", IntegerType), + StructField("s", StringType), + StructField("d", DecimalType(38, 10)), + StructField( + "st", + StructType( + Seq(StructField(nestedName, LongType), StructField("b", ArrayType(StringType))))), + StructField( + "m", + MapType(StringType, ArrayType(StructType(Seq(StructField("x", DoubleType)))))), + StructField("aa", ArrayType(ArrayType(IntegerType))), + StructField("bin", BinaryType), + StructField("ts", TimestampType))) + + private def attributes(nestedName: String = "a"): Seq[AttributeReference] = + schema(nestedName).fields.toSeq.map(f => AttributeReference(f.name, f.dataType, f.nullable)()) + + /** Attributes whose layout no earlier call has cached, since a nested name is part of it. */ + private def newLayout(): Seq[AttributeReference] = + attributes(s"a_${UUID.randomUUID().toString.replace("-", "")}") + + private lazy val rows: Seq[InternalRow] = { + val toCatalyst = CatalystTypeConverters.createToCatalystConverter(schema()) + Seq( + Row( + 1, + "a", + BigDecimal("1.5"), + Row(2L, Seq("x", null)), + Map("k" -> Seq(Row(1.0), null)), + Seq(Seq(1, null), null), + Array[Byte](1, 2), + new Timestamp(0)), + Row(null, null, null, null, null, null, null, null), + Row( + 3, + "字" * 100, + BigDecimal("-1234567890123456789012345678.1234567890"), + Row(null, Seq.empty), + Map.empty, + Seq.empty, + Array.emptyByteArray, + null)).map(toCatalyst(_).asInstanceOf[InternalRow]) + } + + private def withCodegenOnly(confs: (String, String)*)(f: => Unit): Unit = + withSQLConf((SQLConf.CODEGEN_FACTORY_MODE.key -> "CODEGEN_ONLY") +: confs: _*)(f) + + test("rows match UnsafeProjection.create") { + // A small threshold splits the field writes into many methods. + Seq("1024", "64").foreach { threshold => + withCodegenOnly(SQLConf.CODEGEN_METHOD_SPLIT_THRESHOLD.key -> threshold) { + val attrs = newLayout() + val expected = UnsafeProjection.create(attrs, attrs) + val actual = CometUnsafeProjection.create(attrs) + rows.foreach(row => assert(actual(row) == expected(row))) + } + } + } + + test("NO_CODEGEN creates an interpreted projection") { + withSQLConf(SQLConf.CODEGEN_FACTORY_MODE.key -> "NO_CODEGEN") { + val attrs = attributes() + val projection = CometUnsafeProjection.create(attrs) + assert(projection.isInstanceOf[InterpretedUnsafeProjection]) + val expected = UnsafeProjection.create(attrs, attrs) + rows.foreach(row => assert(projection(row) == expected(row))) + } + } + + test("generates the class once per layout and a new projection per call") { + withCodegenOnly() { + val attrs = newLayout() + val before = CometUnsafeProjection.generatedClassCount + val first = CometUnsafeProjection.create(attrs) + val second = CometUnsafeProjection.create(attrs) + assert(CometUnsafeProjection.generatedClassCount == before + 1) + assert(first ne second) + // Each projection writes into its own row, so the first row survives the second call. + val firstRow = first(rows.head) + val secondRow = second(rows(1)) + val expected = UnsafeProjection.create(attrs, attrs) + assert(firstRow == expected(rows.head)) + assert(secondRow == expected(rows(1))) + } + } + + test("a different method split threshold generates its own class") { + val attrs = newLayout() + val before = CometUnsafeProjection.generatedClassCount + Seq("1024", "64", "1024").foreach { threshold => + withCodegenOnly(SQLConf.CODEGEN_METHOD_SPLIT_THRESHOLD.key -> threshold) { + CometUnsafeProjection.create(attrs) + } + } + assert(CometUnsafeProjection.generatedClassCount == before + 2) + } + + test("keeps the most recently used classes") { + withCodegenOnly() { + val layouts = Seq.fill(CometUnsafeProjection.MaxCachedClasses + 1) { + val name = s"f_${UUID.randomUUID().toString.replace("-", "")}" + Seq(AttributeReference("c", StructType(Seq(StructField(name, IntegerType))))()) + } + layouts.foreach(CometUnsafeProjection.create) + val before = CometUnsafeProjection.generatedClassCount + CometUnsafeProjection.create(layouts.last) + assert(CometUnsafeProjection.generatedClassCount == before) + CometUnsafeProjection.create(layouts.head) + assert(CometUnsafeProjection.generatedClassCount == before + 1) + } + } + + test("columnar to row outside whole-stage codegen reuses the generated class") { + withTempPath { dir => + val path = dir.getCanonicalPath + spark + .range(200) + .selectExpr( + "id", + "named_struct('a', id, 'b', array(cast(id as string), null)) as st", + "map(cast(id as string), array(id)) as m") + .repartition(4) + .write + .parquet(path) + // More fields than this keep a plan out of whole-stage codegen, as a wide schema would. + withSQLConf(SQLConf.WHOLESTAGE_MAX_NUM_FIELDS.key -> "2") { + val df = spark.read.parquet(path) + checkSparkAnswerAndOperator(df) + val plan = df.queryExecution.executedPlan + assert(collect(plan) { case c: CometColumnarToRowExec => c }.nonEmpty, plan) + assert(collect(plan) { case w: WholeStageCodegenExec => w }.isEmpty, plan) + // The first run generated the class, so another run over every partition generates none. + val before = CometUnsafeProjection.generatedClassCount + assert(df.collect().length == 200) + assert(CometUnsafeProjection.generatedClassCount == before) + } + } + } +} From 6d3471876f461b0c5ad17c0886dbca4cf9852a75 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Sat, 10 Oct 2026 10:01:32 -0600 Subject: [PATCH 2/2] ci: run CometUnsafeProjectionSuite on Linux and macOS --- .github/workflows/pr_build_linux.yml | 1 + .github/workflows/pr_build_macos.yml | 1 + 2 files changed, 2 insertions(+) diff --git a/.github/workflows/pr_build_linux.yml b/.github/workflows/pr_build_linux.yml index 2bb8fd363d..84a3e39af8 100644 --- a/.github/workflows/pr_build_linux.yml +++ b/.github/workflows/pr_build_linux.yml @@ -560,6 +560,7 @@ jobs: org.apache.spark.sql.comet.CometTPCDSV1_4_PlanStabilitySuite org.apache.spark.sql.comet.CometTPCDSV2_7_PlanStabilitySuite org.apache.spark.sql.comet.CometTaskMetricsSuite + org.apache.spark.sql.comet.CometUnsafeProjectionSuite org.apache.spark.sql.comet.CometDppFallbackRepro3949Suite org.apache.spark.sql.comet.CometShuffleFallbackStickinessSuite org.apache.spark.sql.comet.PlanDataInjectorSuite diff --git a/.github/workflows/pr_build_macos.yml b/.github/workflows/pr_build_macos.yml index 148dde4784..63b95c3bc2 100644 --- a/.github/workflows/pr_build_macos.yml +++ b/.github/workflows/pr_build_macos.yml @@ -264,6 +264,7 @@ jobs: org.apache.spark.sql.comet.CometTPCDSV1_4_PlanStabilitySuite org.apache.spark.sql.comet.CometTPCDSV2_7_PlanStabilitySuite org.apache.spark.sql.comet.CometTaskMetricsSuite + org.apache.spark.sql.comet.CometUnsafeProjectionSuite org.apache.spark.sql.comet.CometDppFallbackRepro3949Suite org.apache.spark.sql.comet.CometShuffleFallbackStickinessSuite org.apache.spark.sql.comet.PlanDataInjectorSuite