unikdahal opened a new issue, #6779:
URL: https://github.com/apache/datafusion-comet/issues/6779
## Description
Comet's broadcast hash join can produce **incorrect results** when the
broadcast side contains variable-width Arrow vectors (string/binary) with
non-zero offsets.
This can occur when the broadcast input is produced by native operators such
as `DISTINCT` / hash aggregation followed by a shuffle.
The issue is reproducible on `main`. Depending on input size, it can either
silently drop matching rows or fail with an Arrow
`OversizedAllocationException`.
## Steps to reproduce
The following regression test reproduces the issue:
```scala
withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10485760",
SQLConf.SHUFFLE_PARTITIONS.key -> "4",
CometConf.COMET_EXEC_ENABLED.key -> "true") {
spark.range(0, 200000, 1, 4)
.selectExpr(
"concat('k', lpad(cast(id as string), 10, '0')) AS id")
.createOrReplaceTempView("t")
checkSparkAnswer(sql("""
SELECT /*+ BROADCAST(d) */ count(*)
FROM t
JOIN (
SELECT DISTINCT
concat('k', lpad(cast(id as string), 10, '0')) AS id
FROM range(0, 200000, 1, 4)
) d ON t.id = d.id
"""))
}
```
**Expected result:** `200000`, matching vanilla Spark.
**Actual result:** Comet returns a different count because some matching
keys are corrupted during broadcast processing.
Inspecting the broadcast batches reveals corrupted string values. For
example, instead of:
```
k0000100008
```
a batch may contain a concatenation of adjacent values:
```
k0000100000k0000100008k0000100016...
```
With larger inputs, the same issue can also result in:
```
OversizedAllocationException:
Memory required for vector is (2147483648)
```
## Suspected root cause
The issue appears to originate in
`spark/src/main/scala/org/apache/spark/sql/comet/util/Utils.scala`,
specifically:
- `serializeBatches`: constructs a `VectorSchemaRoot` using the incoming
Arrow vectors and writes it through `ArrowStreamWriter`.
- `coalesceBroadcastBatches`: merges incoming batches using
`VectorSchemaRootAppender.append`.
Sliced variable-width Arrow vectors can retain non-zero initial offsets
(`offset[0] > 0`) and unused prefixes in their underlying data buffers.
The current broadcast serialization/coalescing path does not explicitly
normalize these offsets. The appender expects offsets relative to the beginning
of the data buffer, causing unused prefixes to become part of the first
appended value.
This corrupts string/binary values and can inflate the computed data size
during repeated appends, eventually triggering oversized allocations.
## Possible fix
Normalize sliced Arrow batches before IPC serialization and broadcast
coalescing.
One approach is to use `VectorSchemaRoot.slice(0, rowCount)`, which uses
Arrow transfer pairs to produce vectors with correctly rebased offsets:
```scala
val normalized = root.slice(0, root.getRowCount)
```
The normalized root can then be passed to `ArrowStreamWriter` or
`VectorSchemaRootAppender.append`.
The implementation should also preserve row counts, handle nested
variable-width vectors, and correctly manage Arrow buffer ownership and cleanup.
Regression coverage should include non-zero-offset string and binary
vectors, serialization round trips, multi-batch coalescing, and an end-to-end
broadcast join with a `DISTINCT` build side.
## Impact
This is a **query correctness issue**
Affected queries can complete successfully while returning incorrect
results, making the problem particularly difficult to detect without comparison
against Spark.
--
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]