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]