comphead commented on code in PR #6371:
URL: https://github.com/apache/datafusion-comet/pull/6371#discussion_r4146857540


##########
spark/src/main/scala/org/apache/spark/sql/comet/util/Utils.scala:
##########
@@ -566,6 +572,68 @@ object Utils extends CometTypeShim with Logging {
     }
   }
 
+  /**
+   * A copy of `large` with 32-bit offsets, as a vector of `narrowType`.
+   *
+   * A PyArrow UDF can hand back a `large_string` or `large_binary` column, 
which Comet reads but
+   * [[getFieldVector]] does not accept. The column is still `StringType` or 
`BinaryType`, which
+   * Comet otherwise carries as `Utf8` or `Binary`, so it is written in that 
form. Whatever reads
+   * the bytes back then gets the type it plans for, and the batches of one 
broadcast share a
+   * schema whichever offset width each arrived with: 
[[coalesceBroadcastBatches]] appends them to
+   * a root built from the first one's schema, and a native plan reading them 
takes its input
+   * stream's schema from the first one too. Only a column holding more data 
than 32-bit offsets
+   * can address is refused.
+   *
+   * Like a materialized `ConstantColumnVector`, the copy is released when the 
caller clears the
+   * vectors it wrote. `large` itself stays with its owner.
+   */
+  private def narrowOffsets(
+      large: BaseLargeVariableWidthVector,
+      narrowType: ArrowType): FieldVector = {
+    val numValues = large.getValueCount
+    val largeOffsets = large.getOffsetBuffer
+    val largeWidth = BaseLargeVariableWidthVector.OFFSET_WIDTH.toLong
+    // An empty vector need not have an offset buffer to read.
+    val start = if (numValues == 0) 0L else largeOffsets.getLong(0)

Review Comment:
   Do the new tests reach the non-zero `start` path? `fill` always writes from 
index 0, so `start` is 0 in all three, and dropping the `- start` below or the 
`start` source index in the data copy would still pass. Would a case like 
`hello`, `null`, `wörld` with offsets `[2, 7, 7, 13]` over `xxhellowörld` be 
worth adding, for example via `loadFieldBuffers`? A zero-row vector also skips 
this code.



##########
spark/src/test/scala/org/apache/spark/sql/comet/util/UtilsSuite.scala:
##########
@@ -197,4 +202,94 @@ class UtilsSuite extends CometTestBase {
       }
     }
   }
+
+  test("serializeBatches writes large-offset string and binary vectors with 
32-bit offsets") {
+    // A PyArrow UDF returning pa.large_string() or pa.large_binary() hands 
Comet a
+    // LargeVarCharVector or LargeVarBinaryVector. The column is still 
StringType or BinaryType,
+    // so it is written as the Utf8 or Binary vector every reader of the bytes 
expects for those.
+    val values = Seq("hello", null, "", "wörld")
+    val numRows = values.length
+    val strings = fill(new LargeVarCharVector("s", CometArrowAllocator), 
values)
+    val binaries = fill(new LargeVarBinaryVector("b", CometArrowAllocator), 
values)
+    try {
+      val batch =
+        new ColumnarBatch(Array(cometVector(strings), cometVector(binaries)), 
numRows)
+
+      val (rowCount, buf) = Utils.serializeBatches(Iterator(batch)).next()
+      assert(rowCount == numRows)
+
+      val it = Utils.decodeBatches(buf, "test")
+      assert(it.hasNext)
+      val out = it.next()
+      val vectorClasses = (0 until out.numCols()).map(i =>
+        out.column(i).asInstanceOf[CometVector].getValueVector.getClass)
+      val gotStrings =
+        (0 until numRows).map(i => 
Option(out.column(0).getUTF8String(i)).map(_.toString).orNull)
+      val gotBinaries = (0 until numRows).map(i =>
+        Option(out.column(1).getBinary(i)).map(new String(_, 
StandardCharsets.UTF_8)).orNull)
+      assert(!it.hasNext)
+
+      assert(vectorClasses == Seq(classOf[VarCharVector], 
classOf[VarBinaryVector]))
+      assert(gotStrings == values)
+      assert(gotBinaries == values)
+    } finally {
+      strings.close()
+      binaries.close()
+    }
+  }
+
+  test("coalesceBroadcastBatches appends a large-offset batch to a 32-bit 
one") {
+    // The coalescer appends every batch to a root built from the first one's 
schema, and leaves
+    // the broadcast uncoalesced when a batch does not match it. A producer 
can hand back
+    // large_string for one batch and string for the next, so both must 
serialize alike.
+    val large = fill(new LargeVarCharVector("s", CometArrowAllocator), 
Seq("a", null))
+    val regular = fill(new VarCharVector("s", CometArrowAllocator), Seq("b", 
"c"))
+    try {
+      val batches = Seq(large, regular).map(v => new 
ColumnarBatch(Array(cometVector(v)), 2))
+      val bufs = 
Utils.serializeBatches(batches.iterator).map(_._2).toSeq.iterator
+      val (coalesced, batchCount, totalRows) = 
Utils.coalesceBroadcastBatches(bufs)
+      // A batch count of zero is the fallback to the uncoalesced buffers.
+      assert(batchCount == 2)
+      assert(totalRows == 4)
+
+      val it = coalesced.iterator.flatMap(Utils.decodeBatches(_, "test"))
+      val out = it.next()
+      val got = (0 until out.numRows()).map(i =>
+        Option(out.column(0).getUTF8String(i)).map(_.toString).orNull)
+      assert(!it.hasNext)
+      assert(got == Seq("a", null, "b", "c"))
+    } finally {
+      large.close()
+      regular.close()
+    }
+  }
+
+  test("serializeBatches refuses a large-offset column too big for 32-bit 
offsets") {
+    // The offsets alone claim 2 GiB of data. The size is checked before 
anything is copied, so
+    // the data itself never has to exist.
+    val strings = fill(new LargeVarCharVector("s", CometArrowAllocator), 
Seq("a"))
+    try {
+      strings.getOffsetBuffer.setLong(
+        BaseLargeVariableWidthVector.OFFSET_WIDTH,
+        Int.MaxValue + 1L)
+      val batch = new ColumnarBatch(Array(cometVector(strings)), 1)
+      val e = 
intercept[SparkException](Utils.serializeBatches(Iterator(batch)).next())
+      assert(e.getMessage.contains("more than 32-bit offsets can address"), 
e.getMessage)
+    } finally {
+      strings.close()
+    }
+  }
+
+  /** Writes `values` into `vector`, leaving a null entry unset so that it 
reads back as null. */
+  private def fill[V <: VariableWidthFieldVector](vector: V, values: 
Seq[String]): V = {

Review Comment:
   Nit: the `isArrowBacked rejects large-offset Arrow vectors` test above 
builds the same `LargeVarCharVector` and `LargeVarBinaryVector` by hand. Would 
it make sense to reuse `fill` and `cometVector` there, so the suite has one way 
to build them?



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