Skip to content

Arrow string-view imports pin whole Parquet pages, so the writer never coalesces string columns #10403

Description

@joseph-isaacs

Summary

String-view columns imported from Arrow keep every Parquet page their views point into, so nbytes() reports up to 15× the bytes the rows use. The writer's coalescing step (RepartitionStrategy, 1 MiB minimum) trusts nbytes(), so every 8,192-row string block looks like it already holds more than 1 MiB and never merges. TPC-H l_comment ends up with 733 chunks of 8,192 rows instead of about 244 chunks of about 24,576 rows.

Cause

  1. arrow-rs's Parquet reader builds each batch's Utf8View/BinaryView views over the whole decompressed page. The TPC-H Parquet stores every string column as Utf8View.
  2. from_arrow_byte_view imports those buffers whole, without copying. Trim child data on Arrow slice imports #10379 trimmed sliced Utf8/Binary/List imports but not view types.
  3. The bench converter's VarBinViewBuilder (compaction_threshold = 0, no deduplication) keeps the buffers and adds a page once for every batch that references it.
  4. VarBinView::slice keeps all data buffers, and nbytes() sums full buffer lengths.

Numeric columns are unaffected because a slice owns exactly len × width bytes.

Evidence (TPC-H SF1, bench converter path)

Column nbytes() fed to writer (B/row) Bytes the rows use (B/row) Chunks written
lineitem.l_comment 637 41 733 × 8,192 rows
orders.o_comment 1,064 65 184 × 8,192 rows
customer.c_comment 1,064 89 19 × 8,192 rows
part.p_name 780 49 25 × 8,192 rows

SpatialBench customer string columns behave the same way. ClickBench and PolarSignals are not affected because their strings import as offset-based Utf8/Binary.

Fix in progress

Branch ji/repartition-view-nbytes, commit d49946a: from_arrow_byte_view now trims view buffers. It drops buffers no valid row references, slices buffers that are at least half used (no copy), and copies sparse values. Fully used buffers are still shared without copying.

Results with the fix:

Before After
l_comment nbytes() fed to writer 637 B/row 46 B/row
l_comment chunks (median rows) 733 (8,192) 244 (24,576)
compress-bench Vortex compress, TPC-H l_comment (whole lineitem) 4.65–4.85 s 3.85–4.18 s (−13 to −20%)
compress-bench Vortex decompress, chunked 211–217 ms 174–194 ms
DataFusion TPC-H Q13 (o_comment), vortex-file-compressed 157 ms 130 ms (−17%)
DataFusion TPC-H Q13, vortex-compact 196 ms 165 ms (−16%)
DataFusion TPC-H, all 22 queries flat (within noise)
string-bench l_comment unchanged

Parquet was used as the control in every run and stayed flat. These numbers come from a 2-core VM with stable 1.97; DuckDB was not measured.

Remaining: what the writer measures

The import fix covers Arrow sources only. Any upstream slice of a VarBinView still makes the coalescer over-measure, because nbytes() reports retained memory.

Proposal: keep nbytes() as retained memory. Change the VarBinView arm of UncompressedSizeInBytes to count views plus referenced bytes, matching the ListView arm, which already rebuilds with MakeExact before measuring. Then have RepartitionStrategy measure blocks with that aggregate instead of nbytes().

Alternatives considered:

  • A second aggregate or a precision option: not needed, because nbytes() is already the cheap upper bound.
  • Canonical::compact() in the coalescer: also unpins memory and moves work the compressor already does, but should_compact's thresholds (50% utilisation, 128 B/row) leave up to 2× slack in the measurement.

The aggregate change should be coordinated with #10177 and #10235.


Generated by Claude Code

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugA bug issue

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions