This is an automated email from the ASF dual-hosted git repository.
Jefffrey pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new 6aac3c2740 fix(parquet): don't preallocate 1MiB in DeltaBitPackEncoder
(#11149)
6aac3c2740 is described below
commit 6aac3c27407f7a44170c1b65cdaceada28a983c5
Author: Peter L <[email protected]>
AuthorDate: Wed Sep 30 18:14:18 2026 +0930
fix(parquet): don't preallocate 1MiB in DeltaBitPackEncoder (#11149)
# Which issue does this PR close?
- Closes #11148.
# Rationale for this change
`DeltaBitPackEncoder::new` preallocates a 1MiB buffer for its bit
writer. Byte array columns eagerly construct their fallback encoder, and
the `DELTA_BYTE_ARRAY` fallback (the default for `PARQUET_2_0`) holds
two `DeltaBitPackEncoder`s, so every byte array column pays 2MiB of heap
up front, even when it dictionary encodes and the fallback is never
used. With 100 columns that is ~200MiB.
The preallocation only ever helped the first page: the buffer is a `Vec`
that retains its capacity across `clear()`.
# What changes are included in this PR?
- `DeltaBitPackEncoder` starts with an empty bit writer buffer that
grows on demand (removes `DEFAULT_BIT_WRITER_SIZE`). This also benefits
`DELTA_BINARY_PACKED` and `DELTA_LENGTH_BYTE_ARRAY`.
- Regression test `unused_delta_fallback_does_not_preallocate` in
`parquet/tests/arrow_writer`, using the existing peak heap tracking
allocator. It fails on `main` (peak ≈ 210MB) and passes with this
change.
# Are these changes tested?
Yes, with the new regression test above. The existing lib and
`arrow_writer` tests pass.
Benchmarks (criterion, `--save-baseline` on `main` vs this branch):
`arrow_writer` delta byte array:
| benchmark | change |
|---|---|
| small_string_shared_prefix/delta_byte_array | −3.8% |
| small_string_partial_prefix/delta_byte_array | −6.6% |
| small_string_distinct/delta_byte_array | −15.6% |
| large_string_shared_prefix/delta_byte_array | −3.7% |
| large_string_distinct/delta_byte_array | −10.4% |
| large_string_shared_prefix_nullable/delta_byte_array | −6.2% |
| large_string_shared_prefix_nullable_dense/delta_byte_array | −15.4% |
| large_string_shared_prefix_nullable_trailing/delta_byte_array | −1.4%
|
| large_string_distinct_nullable/delta_byte_array | −13.0% |
| medium_string_shared_prefix_nullable/delta_byte_array | −3.2% |
| large_string_shared_prefix_list/delta_byte_array | −8.1% |
| string/parquet_2 (dictionary overflow → delta fallback, 3 runs) |
−3.6% to −4.4% |
`writer_overhead`:
| benchmark | change |
|---|---|
| 1000_cols | no change |
| 5000_cols | −2.6% |
| 10000_cols | −2.7% |
| 1000_cols/repeated_batches | −0.6% (noise) |
| 5000_cols/repeated_batches | −14.4% |
The other `parquet_2` / `zstd_parquet_2` writer benchmarks are within
noise after re-running. The one exception is
`int32_ree_95pct_null/parquet_2`, which measured +2.6% to +5.1% across
runs. However, `main` measured against its own baseline drifted +0.9% to
+2.7% on this benchmark, and it does not exercise the delta encoder, so
the difference looks like noise.
# Are there any user-facing changes?
No API changes. Lower peak memory when writing many byte array columns
with `PARQUET_2_0` / delta encodings.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
parquet/src/encodings/encoding/mod.rs | 6 ++++--
parquet/tests/arrow_writer/mod.rs | 29 +++++++++++++++++++++++++++--
2 files changed, 31 insertions(+), 4 deletions(-)
diff --git a/parquet/src/encodings/encoding/mod.rs
b/parquet/src/encodings/encoding/mod.rs
index 4fb8580c70..8b38d4940f 100644
--- a/parquet/src/encodings/encoding/mod.rs
+++ b/parquet/src/encodings/encoding/mod.rs
@@ -370,7 +370,6 @@ impl<T: DataType> Encoder<T> for RleValueEncoder<T> {
// DELTA_BINARY_PACKED encoding
const MAX_PAGE_HEADER_WRITER_SIZE: usize = 32;
-const DEFAULT_BIT_WRITER_SIZE: usize = 1024 * 1024;
const DEFAULT_NUM_MINI_BLOCKS: usize = 4;
/// Delta bit packed encoder.
@@ -434,7 +433,10 @@ impl<T: DataType> DeltaBitPackEncoder<T> {
DeltaBitPackEncoder {
page_header_writer: BitWriter::new(MAX_PAGE_HEADER_WRITER_SIZE),
- bit_writer: BitWriter::new(DEFAULT_BIT_WRITER_SIZE),
+ // Don't pre-allocate: encoders are often created eagerly (e.g. as
a
+ // byte array fallback encoder) and may never be used. The buffer
+ // grows on demand and retains its capacity across pages.
+ bit_writer: BitWriter::new(0),
total_values: 0,
first_value: 0,
current_value: 0, // current value to keep adding deltas
diff --git a/parquet/tests/arrow_writer/mod.rs
b/parquet/tests/arrow_writer/mod.rs
index b1f6a98ddb..c2803c5a13 100644
--- a/parquet/tests/arrow_writer/mod.rs
+++ b/parquet/tests/arrow_writer/mod.rs
@@ -31,7 +31,7 @@ use std::fs::File;
use std::io::{Read as _, Seek, SeekFrom, Write as _};
use std::sync::Arc;
-use arrow::array::{ArrayRef, BinaryArray, Float64Array, Int32Array,
RecordBatch};
+use arrow::array::{ArrayRef, BinaryArray, Float64Array, Int32Array,
RecordBatch, StringArray};
use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
use bytes::Bytes;
use parquet::arrow::arrow_writer::{
@@ -41,7 +41,7 @@ use parquet::arrow::arrow_writer::{
use parquet::arrow::{ArrowSchemaConverter, ArrowWriter};
use parquet::basic::Encoding;
use parquet::errors::Result;
-use parquet::file::properties::WriterProperties;
+use parquet::file::properties::{WriterProperties, WriterVersion};
use parquet::file::writer::SerializedFileWriter;
#[test]
@@ -570,3 +570,28 @@ fn page_store_spills_dictionary_pages() {
(a few dictionary pages), not the ~K × dict_page of the in-memory
baseline"
);
}
+
+/// Regression test for <https://github.com/apache/arrow-rs/issues/11148>
+///
+/// Every byte array column eagerly builds its fallback encoder. With
+/// `DELTA_BYTE_ARRAY` that holds two `DeltaBitPackEncoder`s, which must not
+/// pre-allocate large buffers: here the columns dictionary encode, so the
+/// fallback is never used at all.
+#[test]
+fn unused_delta_fallback_does_not_preallocate() {
+ const COLUMNS: usize = 100;
+ let column: ArrayRef = Arc::new(StringArray::from(vec!["a", "b", "a"]));
+ let batch = RecordBatch::try_from_iter((0..COLUMNS).map(|i|
(format!("c{i}"), column.clone())))
+ .unwrap();
+ let props = WriterProperties::builder()
+ .set_writer_version(WriterVersion::PARQUET_2_0)
+ .set_encoding(Encoding::DELTA_BYTE_ARRAY)
+ .build();
+
+ let peak = peak_heap_bytes(|| {
+ let mut writer = ArrowWriter::try_new(Vec::new(), batch.schema(),
Some(props)).unwrap();
+ writer.write(&batch).unwrap();
+ });
+
+ assert!(peak < COLUMNS * 64 * 1024, "peak: {peak}");
+}