HippoBaro opened a new pull request, #11273:
URL: https://github.com/apache/arrow-rs/pull/11273
> **This is a draft PR and is not intended to be merged as-is.**
Hi all, and @alamb in particular! This PR adds native run-end-encoded (REE)
writes to the Parquet Arrow writer, continuing the run-proportional work
tracked in #9731. It avoids expanding REE inputs into dense arrays, preserving
repetition through nested traversal, level generation, and physical encoding.
### The change
I'm very happy to report that the REE support is now complete! REE can
appear at any level of the schema, and arbitrary compositions such as
`REE<List<REE<Struct<...>>>>` are fully supported. We no longer materialize
dense representations for REE inputs anywhere along the write path.
At a high level, this required redesigning the Arrow-to-Parquet conversion
pipeline. Instead of preparing fully materialized inputs for a collection of
encoding paths, the writer now incrementally describes what to write and lets a
shared, type-aware backend consume the original Arrow storage directly. Once
that underlying architecture was in place, adding REE on top was actually not
all that difficult.
This new architecture also makes almost all workloads faster, particularly
when the old intermediate representation was less memory-efficient than the
original Arrow representation. Boolean arrays are an extreme example: Arrow
gives us a bitmap, but the old path materializes it into a Rust `[bool]`
representation, only to eventually turn it back into a bitset further down the
pipeline. The `bool_non_null/default` workload now reduces to essentially a
memcpy, resulting in ~97% lower runtime (40Ć faster) and ~98% lower peak memory
usage (61Ć less). This is an extreme case, but many other workloads benefit as
well, including `fixed_size_binary*`, `decimal128`, and most dictionary
workloads, to varying degrees (see the performance section below).
### Proposed upstreaming approach
I originally intended to upstream this entirely through small, independently
mergeable PRs, as we started doing with #9653. However, the REE work turned out
to be too complex and wide-ranging for me to confidently publish it
incrementally without risking having to revisit abstractions that had already
been merged. I ultimately had to finish the implementation before I could
settle on a practical upstreaming sequence.
That hopefully explains both the long gap and why Iām now showing up with
exactly the kind of large PR I had hoped to avoid!
I tried to split the work into small, independently mergeable contributions,
but found it difficult to keep every intermediate step both *self-contained and
performance-neutral* without introducing dead code or making individual changes
too large. In many cases, a correct and fully tested intermediate abstraction
adds overhead that only disappears once its consumers are migrated. With the
benchmark suite taking several hours per commit, maintaining all of these
constraints across 30+ commits became completely impractical.
Iād therefore like to separate the **unit of review** from the **unit of
merging**: keep commits small enough to review individually, but group them
into GitHub PR stacks whose complete changes are benchmarked and merged
together.
I hope this is a reasonable compromise š. The 26 commits are organized into
four stacks:
- **A: Establish the contract before changing the machinery.** 9 commits,
ending at `5a628a89c7b1`, covering preliminary correctness fixes and
miscellaneous cleanup. These are mostly independent and can be reviewed and
merged as individual PRs. Several are already published (#11262, #11264,
#11266, #11267, #11268, #11269).
- **B: Teach the encoding pipeline to consume Arrow storage directly.** 8
commits, ending at `d90569926118`, replacing expensive intermediate value
representations with typed encoding families and batch-oriented sources and
sinks.
- **C: Turn structural planning into a stream.** 3 commits, ending at
`d46916b24117`, replacing eager level-and-value planning for ordinary Arrow
trees with record-aligned leaf cursors. This makes structural traversal
incremental while retaining an eager fallback for REE and non-leaf dictionary
compositions.
- **D: Preserve repetition through the entire path.** 6 commits, ending at
`6906adbeafcd`, carrying repeated values as counted selections and adding
native REE traversal, first for scalar leaves and then for nested struct, list,
map, and dictionary compositions.
Every commit on the branch remains a correct, working state with all tests
passing. Performance, however, is guaranteed only at the final commit of each
stack. Stack A contains some breaking bug fixes, while stacks B through D
preserve the existing API.
### Performance
I realize this is a substantial amount of code to work through. Hopefully
the high-level performance results below provide enough of a carrot to make the
review worthwhile š:
```
Compute (`arrow_writer`):
+--------------+----------------+----------------+----------------+
| Baseline | All 483 cases | REE | Non-REE |
+--------------+----------------+----------------+----------------+
| Stack A | baseline | baseline | baseline |
| Stack B | -11.07% | -1.34% | -16.49% |
| Stack C | -17.42% | -2.09% | -25.50% |
| Stack D | -50.46% | -74.55% | -25.87% |
+--------------+----------------+----------------+----------------+
Peak memory (`arrow_writer_peak_memory`):
+--------------+----------------+-------------+-------------+
| Baseline | All 96 cases | REE | Non-REE |
+--------------+----------------+-------------+-------------+
| Stack A | baseline | baseline | baseline |
| Stack B | -21.24% | -9.10% | -29.55% |
| Stack C | -46.02% | -11.09% | -63.38% |
| Stack D | -86.05% | -95.96% | -63.38% |
+--------------+----------------+-------------+-------------+
```
We've also integrated these changes into Datadog's private fork, where we're
seeing roughly **2Ć higher write throughput on production-like, non-synthetic
data**. This isn't directly comparable to the synthetic benchmark suite above,
but it gives us some confidence that the improvements carry over to real-world
workloads.
--
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]