andygrove opened a new pull request, #2557:
URL: https://github.com/apache/datafusion-ballista/pull/2557
# Which issue does this PR close?
Closes #2554.
# Rationale for this change
`ShuffleWriterExec` writes each batch as it comes. For a `Utf8View` or
`BinaryView` column, the IPC writer serializes every data buffer the array
holds, not just the bytes its views point at. TopK, `SortPreservingMergeExec`
and `LIMIT` emit a few rows that still hold the data buffers of the batches
those rows came from, so a tiny result can be written as tens of MB. With the
reproduction in #2554, Q10's 20-row result is a 26 MB file, and the client
can't fetch it at the default 16 MiB gRPC limit.
The sort-based shuffle writer already compacts view columns before it writes
them (`compact_view_columns` in `partitioned_batch_iterator.rs`). The
passthrough writers don't.
# What changes are included in this PR?
- A new `compact_sparse_view_columns` helper in `utils.rs`. It runs `gc()`
on a top-level `Utf8View` / `BinaryView` column when the column's data buffers
hold more than twice the bytes its views reference. That's the threshold
arrow's `BatchCoalescer` uses, and the one #2521 proposes. Dense columns are
written as they are, so the copy only happens when it saves at least half the
bytes.
- `write_stream_to_disk` (used by `ShuffleWriterExec`) and
`write_stream_to_ipc_file` (used by `RangeShuffleWriterExec`, which replaces a
passthrough writer on sorted, range-partitioned stages) call the helper on each
batch before writing. It runs in their existing blocking writer task, so the
copy stays off the async runtime and counts toward `write_time`.
- The sort-shuffle writer still compacts unconditionally. Moving it onto
this helper is what #2521 asks for, so I left it out. Neither writer compacts
view columns nested inside a struct or list.
Two unit tests cover the helper: sparse `Utf8View` and `BinaryView` columns
are compacted, and a dense column comes back as the same array. A regression
test for each writer writes 2 rows sliced from 1,000 long strings and checks
that the file holds only those 2 rows' bytes. Both writer tests fail on `main`,
with 32,000 bytes written against 64 referenced.
End to end, I used the reproduction from #2554: TPC-H SF10 from
`tpchgen-cli`, with `customer` rewritten by datafusion-cli 53 so its strings
decode as `Utf8View`. The cluster was one scheduler and one executor with
default gRPC limits.
| | `main`
| this PR |
| -------------- |
--------------------------------------------------------------- | ------------:
|
| Q10 | fails, `encoded message length too large: found 26758377
bytes` | 20 rows |
| Stage 4 output | 58,062,272 bytes
| 75,392 bytes |
| Stage 5 output | 26,760,584 bytes
| 5,640 bytes |
The unmodified `tpchgen-cli` files give 69,760 and 5,000 bytes for those
stages, and the same 20 rows.
I also ran all 22 TPC-H queries on three copies of SF10: the `tpchgen-cli`
files, the copy above, and a copy where datafusion-cli rewrote `lineitem`,
`orders`, `customer` and `part` as Hive-partitioned tables. `main` fails Q10 on
both rewritten copies. With this PR all 22 queries pass on all three, with the
same row counts.
# Are there any user-facing changes?
No API changes. Queries over `Utf8View` / `BinaryView` data write smaller
shuffle files, and small results that used to exceed the gRPC message limit now
fit.
#2165 reports the same error on Q21 and Q22. It may have the same cause, but
the report doesn't say what data it ran on, so this PR doesn't claim to fix it.
--
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]