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


##########
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:
   This starts well below the 64 KB limit. Under it the method still compiles, 
but once `next()` passes HotSpot's 8000-byte `HugeMethodLimit` the JIT never 
compiles it and it stays interpreted. I timed this iterator against main's 
`UnsafeProjection` plus `copy()` over nullable bigint columns, 8 batches of 
4096 rows each. It is 0.39x at 6 columns and about the same as main at 50. It 
is 1.7x slower at 100 columns and 13 to 15x slower from 120 to 500. String 
columns reach 4x at 150. With `-XX:-DontCompileHugeMethods` the 120 and 150 
column cases drop back to about 1.2x, which points at the method size. End to 
end, reading a 150-column Comet cache took 320 ms on this PR against 66 ms on 
the merge base with Comet on and exec off, and 323 ms against 64 ms with Comet 
off. Spark reads of relations wider than `spark.sql.codegen.maxFields` always 
take this path, because `InMemoryTableScanExec.supportsColumnar` is false for 
them.
   
   So falling back only when compilation fails would not be enough. 
`CodeGenerator.compile` returns the `ByteCodeStats` that line 110 discards, and 
that is what `WholeStageCodegenExec` checks before it backs off. A width bound 
would work too. For the fallback itself, `UnsafeProjection.create(fields)` over 
`batch.getRow(i)` in an indexed loop, without the `copy()`, came in at 0.64 to 
0.93x of main at every width I tried. Could the benchmark and the tests cover 
100 and 200 columns as well as 1,500?



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