Repository navigation
perf: generate the columnar-to-row projection once per executor instead of per partition #6858
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
andygrove
wants to merge
2
commits into
apache:main
Choose a base branch
from
andygrove:andygrove/datafusion-comet-6855
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
Show all changes
2 commits
Select commit
Hold shift + click to select a range
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
146 changes: 146 additions & 0 deletions
146
spark/src/main/scala/org/apache/spark/sql/comet/CometUnsafeProjection.scala
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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) | ||
| } | ||
| } |
177 changes: 177 additions & 0 deletions
177
spark/src/test/scala/org/apache/spark/sql/comet/CometUnsafeProjectionSuite.scala
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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 { | ||
|
|
||
| /** 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) | ||
| } | ||
| } | ||
| } | ||
| } | ||
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
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
CometUnsafeProjectionSuitewithout listing its fully qualified name in.github/workflows/pr_build_linux.ymland.github/workflows/pr_build_macos.ymlmakes 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.pyexits 255 withSuite 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, andRequired Checksfails.There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Registered
CometUnsafeProjectionSuitein the exec group of both workflows in 6d34718;check-suites.pyand Preflight pass.