Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -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)
}
}
Original file line number Diff line number Diff line change
@@ -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 {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Register the new suite in both CI workflow matrices. Adding CometUnsafeProjectionSuite without listing its fully qualified name in .github/workflows/pr_build_linux.yml and .github/workflows/pr_build_macos.yml makes the mandatory suite-inventory check fail. Expected behavior is for preflight to pass and the new regression tests to run. Instead, exact-head preflight fails and downstream builds/tests are skipped. Please add the suite to the appropriate group in both workflows.

Evidence: At ed8ba0c, python3 dev/ci/check-suites.py exits 255 with Suite not found in workflow .github/workflows/pr_build_linux.yml: org.apache.spark.sql.comet.CometUnsafeProjectionSuite. The suite name is absent from both workflow files. Exact-head CI reproduces this at https://github.com/apache/datafusion-comet/actions/runs/38063389023/job/114246025282, and Required Checks fails.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Registered CometUnsafeProjectionSuite in the exec group of both workflows in 6d34718; check-suites.py and Preflight pass.


/** 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)
}
}
}
}
Loading