andygrove opened a new pull request, #6371:
URL: https://github.com/apache/datafusion-comet/pull/6371

   ## Which issue does this PR close?
   
   Closes #5488.
   
   ## Rationale for this change
   
   `Utils.getFieldVector` accepts a fixed list of Arrow vector classes that 
leaves out `LargeVarCharVector` and `LargeVarBinaryVector`, so 
`Utils.serializeBatches` throws `Unsupported Arrow Vector for serialize` on a 
batch that holds either. Comet reads these vectors on purpose elsewhere. A 
PyArrow UDF can return `large_string` or `large_binary` (pandas 3 backs string 
columns with them), `CometPlainVector` reads 64-bit offsets, and 
`Utils.toArrowType` maps `LargeUtf8` to `StringType`.
   
   The gap is latent on `main`, and I have not reproduced it with a query. 
`serializeBatches` has two callers. `getByteArrayRdd` in `operators.scala` has 
no callers. `CometBroadcastExchangeExec` only broadcasts a plan whose children 
are all `CometNativeExec`. `CometMapInBatchExec`, the operator that keeps the 
Python worker's vectors, is not one of those. `EliminateRedundantTransitions` 
also inserts it after `CometExecRule` has run, so its output reaches the rest 
of the plan through a row transition. This change makes the serialization 
boundary hold once such output can reach it.
   
   The issue suggested two fixes: accept the large vectors in `getFieldVector`, 
or narrow their offsets to 32 bits before serializing. This PR narrows them, 
because of what happens to the bytes after `serializeBatches`:
   
   - On the build side of a broadcast join, `CometBatchRDD` decodes the bytes 
with `Utils.decodeBatches`, and `CometArrowStream.wrapColumnarBatchRDD` feeds 
the batches to the native plan. `reconcileStreamSchema` advertises the first 
batch's actual Arrow types for the whole stream, and `ColumnarBatchArrowReader` 
loads every later batch's buffers under that schema. If each batch kept its own 
offset width, a broadcast that holds both widths would load one width's offset 
buffers under the other's type.
   - `Utils.coalesceBroadcastBatches` appends every batch into a root built 
from the first batch's schema. When a later batch doesn't match, it leaves the 
broadcast uncoalesced, so mixed widths would also lose coalescing.
   - `CometNativeColumnarToRowExec` hands decoded broadcast batches to native 
code through `NativeUtil.exportBatch`, which checks each vector with the same 
`getFieldVector`. Accepting large vectors there would widen what every 
JVM-to-native export accepts, not just serialization.
   - The column's Spark type is `StringType` or `BinaryType`, which Comet 
carries as `Utf8` or `Binary` everywhere else, including what the cache 
serializer writes when it converts a batch. Writing that form gives every 
reader the type it plans for. Without narrowing, native `ScanExec` would cast 
`LargeUtf8` back to the planned `Utf8` in every task that reads the broadcast. 
Narrowing copies the column once, when it is serialized.
   
   ## What changes are included in this PR?
   
   - `Utils.getBatchFieldVectorsWithProviders` builds the vectors 
`serializeBatches` writes. It now replaces a `LargeVarCharVector` or 
`LargeVarBinaryVector` with a `VarCharVector` or `VarBinaryVector` copy, made 
by a new private `Utils.narrowOffsets`.
     - The copy moves the validity and data buffers in bulk, rewrites the 
offsets, and keeps the field's name, nullability and metadata.
     - A column whose data does not fit 32-bit offsets (more than 2 GiB in one 
batch) fails with a `SparkException` that names the column. The size is checked 
before anything is allocated.
     - Like a materialized `ConstantColumnVector` in the same method, the copy 
is released when `serializeBatches` clears its root. The original vector stays 
with its owner.
   - `getFieldVector` and `isSupportedFieldVector` are unchanged and still 
agree. The FFI export in `NativeUtil.exportBatch` still rejects large-offset 
vectors. `isArrowBacked` still answers false for them, so the cache serializer 
keeps converting such a batch. The comments explaining that answer are updated, 
since writing such a batch no longer fails.
   - Nested large-offset children, such as a struct of `large_string`, already 
pass `getFieldVector` and are written as they are. This PR leaves them alone.
   
   ## How are these changes tested?
   
   New tests in `UtilsSuite`:
   
   - `serializeBatches writes large-offset string and binary vectors with 
32-bit offsets`: a batch with a `LargeVarCharVector` and a 
`LargeVarBinaryVector` column, whose values include a null, an empty value and 
a multi-byte character, is serialized and then decoded with 
`Utils.decodeBatches`, as the broadcast consumers do. The decoded columns are 
`VarCharVector` and `VarBinaryVector` and hold the original values.
   - `coalesceBroadcastBatches appends a large-offset batch to a 32-bit one`: a 
large-offset batch followed by a `VarCharVector` batch coalesces into a single 
batch, without falling back, and keeps all four values in order.
   - `serializeBatches refuses a large-offset column too big for 32-bit 
offsets`: offsets that claim more than `Int.MaxValue` bytes fail with the new 
message, before anything is copied.
   
   Without the change to `Utils.scala`, all three fail with `SparkException: 
Unsupported Arrow Vector for serialize: class 
org.apache.arrow.vector.LargeVarCharVector`.
   
   With the change, on the default profile (Spark 4.1, Scala 2.13), I ran 
`UtilsSuite`, `NativeUtilSuite`, `CometInMemoryCacheSuite` and the 
`CometJoinSuite` tests whose names contain `Broadcast`. All 94 tests pass.
   


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