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]

Reply via email to