sunchao commented on code in PR #5859:
URL: https://github.com/apache/datafusion-comet/pull/5859#discussion_r4104531486


##########
spark/src/main/scala/org/apache/spark/sql/comet/execution/arrow/CachedBatchRowIterator.scala:
##########
@@ -0,0 +1,131 @@
+/*
+ * 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.execution.arrow
+
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions.{Attribute, BoundReference, 
CodeGeneratorWithInterpretedFallback, InterpretedUnsafeProjection}
+import org.apache.spark.sql.catalyst.expressions.codegen._
+import org.apache.spark.sql.catalyst.expressions.codegen.Block._
+import org.apache.spark.sql.vectorized.{ColumnarBatch, ColumnVector}
+
+/**
+ * Reads vectors directly into Spark's reusable UnsafeRow buffer. The input 
iterator owns the
+ * batches and releases them on advancement or task completion. As with 
Spark's cache reader,
+ * callers must copy rows they retain across next(), but the returned row owns 
its variable-width
+ * values and remains valid when hasNext() releases the batch that supplied 
them.
+ */
+private[arrow] class CachedBatchRowIterator(attributes: Seq[Attribute])
+    extends CodeGeneratorWithInterpretedFallback[Iterator[ColumnarBatch], 
Iterator[InternalRow]] {
+
+  private def fields: Seq[BoundReference] = attributes.zipWithIndex.map { case 
(attr, i) =>
+    BoundReference(i, attr.dataType, attr.nullable)
+  }
+
+  override protected def createCodeGeneratedObject(
+      batches: Iterator[ColumnarBatch]): Iterator[InternalRow] = {
+    val ctx = new CodegenContext
+    val columns = attributes.indices.map { i =>
+      ctx.addMutableState(classOf[ColumnVector].getName, s"column$i")
+    }
+    ctx.currentVars = attributes.zip(columns).map { case (attr, column) =>
+      val value = JavaCode.variable(ctx.freshName("value"), attr.dataType)
+      val getter = CodeGenerator.getValueFromVector(column, attr.dataType, 
"rowId")
+      val javaType = CodeGenerator.javaType(attr.dataType)
+      if (attr.nullable) {
+        val isNull = JavaCode.isNullVariable(ctx.freshName("isNull"))
+        ExprCode(
+          code"""
+            boolean $isNull = $column.isNullAt(rowId);
+            $javaType $value = $isNull ? 
${CodeGenerator.defaultValue(attr.dataType)} : ($getter);
+          """,
+          isNull,
+          value)
+      } else {
+        ExprCode(code"$javaType $value = $getter;", FalseLiteral, value)
+      }
+    }
+    val projection = GenerateUnsafeProjection.createCode(ctx, fields)

Review Comment:
   [P2] Preserve code splitting for wide cache projections. Reading 1,500 
selected INT columns with alternating nullable/non-nullable attributes makes 
the generated `next()` exceed the JVM's 64 KB method limit. Setting 
`ctx.currentVars` above causes `GenerateUnsafeProjection` to skip its normal 
`splitExpressions` path, whereas the previous `UnsafeProjection.create(...)` 
reader compiles successfully for the same schema. Under default `FALLBACK`, 
each partition retries the failed compilation and then uses the interpreted 
reader. The local probe measured later constructions at 400–515 ms versus 43–63 
ms for the previous path, with correct results from both. `CODEGEN_ONLY` fails 
outright. Could this split the generated writer into bounded methods, or select 
the existing generated `UnsafeProjection` path for wide schemas before 
attempting oversized compilation? A 1,500-column regression case would cover 
this.
   
   Evidence: Compiled the unchanged exact-head `CachedBatchRowIterator.scala` 
against Spark 4.1.3 with Scala 2.13.17/JDK 21. A package-local harness created 
attributes using `(0 until 1500).map(i => AttributeReference(s"c$i", 
IntegerType, nullable = i % 2 == 0)())` and matching `OnHeapColumnVector`s. 
With `CODEGEN_ONLY`, the previous `UnsafeProjection.create(attrs, attrs)` path 
passed, while `new 
CachedBatchRowIterator(attrs).createObject(Iterator.single(batch))` failed with 
`InternalCompilerException: Code grows beyond 64 KB` while compiling `next()`. 
Widths 150, 500 and 1,000 passed. In a separate default-`FALLBACK` probe with 
4,096 rows and logging disabled, five alternating old/new runs returned 
identical checksums. The final three construction times were 62.95/45.58/43.28 
ms for the previous path and 399.82/484.91/514.83 ms for the new interpreted 
fallback. Spark's `GenerateUnsafeProjection.writeExpressionsToBuffer` 
explicitly bypasses splitting when `ctx.currentVars != null`. H
 arnesses and logs are under `/tmp/comet-5859-dbfb-probe/`.



-- 
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]

Reply via email to