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]