From 490b2a697f935c80cb9cb2b689536ec1b1d4c267 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Fri, 18 Sep 2026 11:49:49 +0200 Subject: [PATCH 1/3] perf(parquet): reuse delta discard storage --- parquet/internal/encoding/delta_byte_array.go | 37 ++++++++--- .../delta_byte_array_decode_benchmark_test.go | 53 ++++++++++++++++ .../encoding/delta_byte_array_decode_test.go | 63 +++++++++++++++++++ 3 files changed, 144 insertions(+), 9 deletions(-) diff --git a/parquet/internal/encoding/delta_byte_array.go b/parquet/internal/encoding/delta_byte_array.go index b2ed5bf53..ca97341e8 100644 --- a/parquet/internal/encoding/delta_byte_array.go +++ b/parquet/internal/encoding/delta_byte_array.go @@ -19,6 +19,7 @@ package encoding import ( "errors" "fmt" + "slices" "github.com/apache/arrow-go/v18/arrow/memory" "github.com/apache/arrow-go/v18/internal/utils" @@ -162,8 +163,9 @@ func (enc *DeltaByteArrayEncoder) FlushValues() (Buffer, error) { type DeltaByteArrayDecoder struct { *DeltaLengthByteArrayDecoder - prefixLengths []int32 - lastVal parquet.ByteArray + prefixLengths []int32 + lastVal parquet.ByteArray + discardScratch []byte } // Type returns the underlying physical type this decoder operates on, in this case ByteArrays only @@ -173,6 +175,25 @@ func (DeltaByteArrayDecoder) Type() parquet.Type { func (d *DeltaByteArrayDecoder) Allocator() memory.Allocator { return d.mem } +func (d *DeltaByteArrayDecoder) setDiscardLastValue(prefix, suffix parquet.ByteArray) { + valueLen := len(prefix) + len(suffix) + if valueLen == 0 { + if d.discardScratch == nil { + d.discardScratch = make([]byte, 0, 1) + } else { + d.discardScratch = d.discardScratch[:0] + } + d.lastVal = d.discardScratch + return + } + + d.discardScratch = slices.Grow(d.discardScratch[:0], valueLen) + d.discardScratch = d.discardScratch[:valueLen] + copy(d.discardScratch, prefix) + copy(d.discardScratch[len(prefix):], suffix) + d.lastVal = d.discardScratch +} + // SetData expects the passed in data to be the prefix lengths, followed by the // blocks of suffix data in order to initialize the decoder. func (d *DeltaByteArrayDecoder) SetData(nvalues int, data []byte) error { @@ -222,15 +243,15 @@ func (d *DeltaByteArrayDecoder) Discard(n int) (int, error) { } remaining := n - tmp := make([]parquet.ByteArray, 1) + var tmp [1]parquet.ByteArray if d.lastVal == nil { if len(d.prefixLengths) == 0 || d.prefixLengths[0] != 0 { return 0, errors.New("parquet: first delta byte array prefix length must be zero") } - if _, err := d.DeltaLengthByteArrayDecoder.Decode(tmp); err != nil { + if _, err := d.DeltaLengthByteArrayDecoder.Decode(tmp[:]); err != nil { return 0, err } - d.lastVal = tmp[0] + d.setDiscardLastValue(nil, tmp[0]) d.prefixLengths = d.prefixLengths[1:] remaining-- } @@ -246,16 +267,14 @@ func (d *DeltaByteArrayDecoder) Discard(n int) (int, error) { } prefix := d.lastVal[:prefixLen:prefixLen] - if _, err := d.DeltaLengthByteArrayDecoder.Decode(tmp); err != nil { + if _, err := d.DeltaLengthByteArrayDecoder.Decode(tmp[:]); err != nil { return n - remaining, err } if len(tmp[0]) == 0 { d.lastVal = prefix } else { - d.lastVal = make([]byte, int(prefixLen)+len(tmp[0])) - copy(d.lastVal, prefix) - copy(d.lastVal[prefixLen:], tmp[0]) + d.setDiscardLastValue(prefix, tmp[0]) } remaining-- } diff --git a/parquet/internal/encoding/delta_byte_array_decode_benchmark_test.go b/parquet/internal/encoding/delta_byte_array_decode_benchmark_test.go index b03df3a03..0ff9a936c 100644 --- a/parquet/internal/encoding/delta_byte_array_decode_benchmark_test.go +++ b/parquet/internal/encoding/delta_byte_array_decode_benchmark_test.go @@ -95,6 +95,59 @@ func BenchmarkDeltaByteArrayDecoderDecode(b *testing.B) { } } +func BenchmarkDeltaByteArrayDecoderDiscard(b *testing.B) { + for _, test := range []struct { + name string + value func(int) string + }{ + { + name: "prefix-heavy", + value: func(i int) string { + return fmt.Sprintf("tenant/%04d/partition/%04d/object", i/100, i) + }, + }, + { + name: "low-prefix", + value: func(i int) string { + return fmt.Sprintf("%08x/%08x", i, i*7919) + }, + }, + } { + for _, nvalues := range []int{1024, 65536} { + test := test + nvalues := nvalues + b.Run(fmt.Sprintf("%s/%d", test.name, nvalues), func(b *testing.B) { + values := make([]parquet.ByteArray, nvalues) + inputBytes := 0 + for i := range values { + values[i] = parquet.ByteArray(test.value(i)) + inputBytes += len(values[i]) + } + encoded := encodeDeltaByteArrayValues(values) + dec := NewDecoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray, + nil, memory.DefaultAllocator).(*DeltaByteArrayDecoder) + b.SetBytes(int64(inputBytes)) + b.ReportAllocs() + b.ResetTimer() + for b.Loop() { + b.StopTimer() + if err := dec.SetData(nvalues, encoded); err != nil { + b.Fatal(err) + } + b.StartTimer() + discarded, err := dec.Discard(nvalues) + if err != nil { + b.Fatal(err) + } + if discarded != nvalues { + b.Fatalf("discarded %d values, expected %d", discarded, nvalues) + } + } + }) + } + } +} + func encodeDeltaByteArrayValues(values []parquet.ByteArray) []byte { enc := NewEncoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray, false, nil, memory.DefaultAllocator).(ByteArrayEncoder) diff --git a/parquet/internal/encoding/delta_byte_array_decode_test.go b/parquet/internal/encoding/delta_byte_array_decode_test.go index 1a5ea5f25..8a4f0bf4f 100644 --- a/parquet/internal/encoding/delta_byte_array_decode_test.go +++ b/parquet/internal/encoding/delta_byte_array_decode_test.go @@ -17,6 +17,7 @@ package encoding import ( + "bytes" "fmt" "strings" "testing" @@ -92,6 +93,68 @@ func TestDeltaByteArrayDecoderDecodesAllEmptyValues(t *testing.T) { } } +func TestDeltaByteArrayDecoderDiscardsAllEmptyValues(t *testing.T) { + values := []string{"", "", ""} + data := encodeDeltaByteArrayPage(t, values) + dec := NewDecoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray, + nil, memory.DefaultAllocator).(ByteArrayDecoder) + require.NoError(t, dec.SetData(len(values), data)) + + discarded, err := dec.Discard(1) + require.NoError(t, err) + require.Equal(t, 1, discarded) + + discarded, err = dec.Discard(2) + require.NoError(t, err) + require.Equal(t, 2, discarded) +} + +func TestDeltaByteArrayDecoderDiscardCopiesFirstValue(t *testing.T) { + values := []string{"first-value", "first-value", "first-value/final"} + data := encodeDeltaByteArrayPage(t, values) + dec := NewDecoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray, + nil, memory.DefaultAllocator).(ByteArrayDecoder) + require.NoError(t, dec.SetData(len(values), data)) + + discarded, err := dec.Discard(2) + require.NoError(t, err) + require.Equal(t, 2, discarded) + + firstValueOffset := bytes.Index(data, []byte(values[0])) + require.NotEqual(t, -1, firstValueOffset) + copy(data[firstValueOffset:firstValueOffset+len(values[0])], strings.Repeat("x", len(values[0]))) + + out := make([]parquet.ByteArray, 1) + decoded, err := dec.Decode(out) + require.NoError(t, err) + require.Equal(t, 1, decoded) + require.Equal(t, values[2], string(out[0])) +} + +func TestDeltaByteArrayDecoderReusesDiscardScratch(t *testing.T) { + firstValues := []string{"prefix/000", "prefix/001", "prefix/002", "prefix/003"} + firstData := encodeDeltaByteArrayPage(t, firstValues) + dec := NewDecoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray, + nil, memory.DefaultAllocator).(*DeltaByteArrayDecoder) + require.NoError(t, dec.SetData(len(firstValues), firstData)) + + discarded, err := dec.Discard(len(firstValues)) + require.NoError(t, err) + require.Equal(t, len(firstValues), discarded) + require.NotEmpty(t, dec.discardScratch) + scratchStart := &dec.discardScratch[0] + scratchCap := cap(dec.discardScratch) + + secondValues := []string{"prefix/100", "prefix/101"} + secondData := encodeDeltaByteArrayPage(t, secondValues) + require.NoError(t, dec.SetData(len(secondValues), secondData)) + discarded, err = dec.Discard(len(secondValues)) + require.NoError(t, err) + require.Equal(t, len(secondValues), discarded) + require.Equal(t, scratchCap, cap(dec.discardScratch)) + require.Equal(t, scratchStart, &dec.discardScratch[0]) +} + func TestDeltaByteArrayDecoderReusesValuesWithoutSuffixes(t *testing.T) { value := strings.Repeat("x", 64*1024) values := make([]string, 128) From 2956ddf2516ef1508466137985df0b7eda06c5eb Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Mon, 21 Sep 2026 07:29:42 +0200 Subject: [PATCH 2/3] fix(parquet): keep decoded prefixes out of discard scratch --- parquet/internal/encoding/delta_byte_array.go | 4 + .../delta_byte_array_discard_lifetime_test.go | 92 +++++++++++++++++++ 2 files changed, 96 insertions(+) create mode 100644 parquet/internal/encoding/delta_byte_array_discard_lifetime_test.go diff --git a/parquet/internal/encoding/delta_byte_array.go b/parquet/internal/encoding/delta_byte_array.go index ca97341e8..d8f76dbc0 100644 --- a/parquet/internal/encoding/delta_byte_array.go +++ b/parquet/internal/encoding/delta_byte_array.go @@ -382,6 +382,10 @@ func (d *DeltaByteArrayDecoder) Decode(out []parquet.ByteArray) (int, error) { prefix := d.lastVal[:prefixLen:prefixLen] if len(out[0]) == 0 { + // Decoded values must not escape through reusable discard storage. + if len(prefix) > 0 && len(d.discardScratch) > 0 && &prefix[0] == &d.discardScratch[0] { + prefix = slices.Clone(prefix) + } d.lastVal = prefix out[0], out = prefix, out[1:] continue diff --git a/parquet/internal/encoding/delta_byte_array_discard_lifetime_test.go b/parquet/internal/encoding/delta_byte_array_discard_lifetime_test.go new file mode 100644 index 000000000..f9f289dee --- /dev/null +++ b/parquet/internal/encoding/delta_byte_array_discard_lifetime_test.go @@ -0,0 +1,92 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package encoding + +import ( + "testing" + + "github.com/apache/arrow-go/v18/arrow/memory" + "github.com/apache/arrow-go/v18/parquet" + "github.com/stretchr/testify/require" +) + +func TestDeltaByteArrayDecoderDiscardKeepsDecodedPrefixes(t *testing.T) { + for _, tc := range []struct { + name string + value string + }{ + {"repeated", "aa"}, + {"shorter-prefix", "a"}, + {"empty", ""}, + {"non-empty-suffix", "ab"}, + } { + for _, nextPage := range []bool{false, true} { + pageName := "same-page" + if nextPage { + pageName = "next-page" + } + t.Run(tc.name+"/"+pageName, func(t *testing.T) { + values := []string{"aa", tc.value, "zz"} + if nextPage { + values = values[:2] + } + dec := NewDecoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray, + nil, memory.DefaultAllocator).(ByteArrayDecoder) + require.NoError(t, dec.SetData(len(values), encodeDeltaByteArrayPage(t, values))) + + discarded, err := dec.Discard(1) + require.NoError(t, err) + require.Equal(t, 1, discarded) + + out := make([]parquet.ByteArray, 1) + decoded, err := dec.Decode(out) + require.NoError(t, err) + require.Equal(t, 1, decoded) + require.Equal(t, tc.value, string(out[0])) + + if nextPage { + require.NoError(t, dec.SetData(1, encodeDeltaByteArrayPage(t, []string{"zz"}))) + } + discarded, err = dec.Discard(1) + require.NoError(t, err) + require.Equal(t, 1, discarded) + require.Equal(t, tc.value, string(out[0])) + }) + } + } +} + +func TestDeltaByteArrayDecoderDiscardKeepsDecodedPrefixBatch(t *testing.T) { + values := []string{"aa", "aa", "a", "zz"} + dec := NewDecoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray, + nil, memory.DefaultAllocator).(ByteArrayDecoder) + require.NoError(t, dec.SetData(len(values), encodeDeltaByteArrayPage(t, values))) + + discarded, err := dec.Discard(1) + require.NoError(t, err) + require.Equal(t, 1, discarded) + out := make([]parquet.ByteArray, 2) + decoded, err := dec.Decode(out) + require.NoError(t, err) + require.Equal(t, 2, decoded) + + discarded, err = dec.Discard(1) + require.NoError(t, err) + require.Equal(t, 1, discarded) + require.Equal(t, "aa", string(out[0])) + require.Equal(t, "a", string(out[1])) +} From 6b9fbf9f28e96acb2a21c8dd2e34ec8847708421 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Mon, 21 Sep 2026 09:51:19 +0200 Subject: [PATCH 3/3] perf(parquet): read discard suffixes without single-element batches --- parquet/internal/encoding/delta_byte_array.go | 27 ++++--- .../delta_byte_array_discard_chunks_test.go | 72 +++++++++++++++++++ 2 files changed, 88 insertions(+), 11 deletions(-) create mode 100644 parquet/internal/encoding/delta_byte_array_discard_chunks_test.go diff --git a/parquet/internal/encoding/delta_byte_array.go b/parquet/internal/encoding/delta_byte_array.go index d8f76dbc0..d55ddaa1d 100644 --- a/parquet/internal/encoding/delta_byte_array.go +++ b/parquet/internal/encoding/delta_byte_array.go @@ -243,15 +243,12 @@ func (d *DeltaByteArrayDecoder) Discard(n int) (int, error) { } remaining := n - var tmp [1]parquet.ByteArray if d.lastVal == nil { if len(d.prefixLengths) == 0 || d.prefixLengths[0] != 0 { return 0, errors.New("parquet: first delta byte array prefix length must be zero") } - if _, err := d.DeltaLengthByteArrayDecoder.Decode(tmp[:]); err != nil { - return 0, err - } - d.setDiscardLastValue(nil, tmp[0]) + suffix := d.decodeDiscardSuffix() + d.setDiscardLastValue(nil, suffix) d.prefixLengths = d.prefixLengths[1:] remaining-- } @@ -267,14 +264,11 @@ func (d *DeltaByteArrayDecoder) Discard(n int) (int, error) { } prefix := d.lastVal[:prefixLen:prefixLen] - if _, err := d.DeltaLengthByteArrayDecoder.Decode(tmp[:]); err != nil { - return n - remaining, err - } - - if len(tmp[0]) == 0 { + suffix := d.decodeDiscardSuffix() + if len(suffix) == 0 { d.lastVal = prefix } else { - d.setDiscardLastValue(prefix, tmp[0]) + d.setDiscardLastValue(prefix, suffix) } remaining-- } @@ -282,6 +276,17 @@ func (d *DeltaByteArrayDecoder) Discard(n int) (int, error) { return n, nil } +// decodeDiscardSuffix reads one suffix after Discard has bounded its count by +// nvals. SetData has already validated the suffix lengths against the payload. +func (d *DeltaByteArrayDecoder) decodeDiscardSuffix() parquet.ByteArray { + length := d.lengths[0] + suffix := d.data[:length:length] + d.data = d.data[length:] + d.nvals-- + d.lengths = d.lengths[1:] + return suffix +} + func (d *DeltaByteArrayDecoder) decodedArenaSize(max int) (int, error) { maxInt := int(^uint(0) >> 1) total := 0 diff --git a/parquet/internal/encoding/delta_byte_array_discard_chunks_test.go b/parquet/internal/encoding/delta_byte_array_discard_chunks_test.go new file mode 100644 index 000000000..718a16ed6 --- /dev/null +++ b/parquet/internal/encoding/delta_byte_array_discard_chunks_test.go @@ -0,0 +1,72 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package encoding + +import ( + "fmt" + "testing" + + "github.com/apache/arrow-go/v18/arrow/memory" + "github.com/apache/arrow-go/v18/parquet" + "github.com/stretchr/testify/require" +) + +func TestDeltaByteArrayDecoderDiscardChunkBoundaries(t *testing.T) { + values := []string{"aa", "aa", "a", "", "", "prefix/000", "prefix/001", "prefix/001", "z"} + data := encodeDeltaByteArrayPage(t, values) + for initial := 0; initial <= len(values); initial++ { + for skip := 0; skip <= len(values)+1; skip++ { + t.Run(fmt.Sprintf("decoded=%d/discard=%d", initial, skip), func(t *testing.T) { + dec := NewDecoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray, + nil, memory.DefaultAllocator).(*DeltaByteArrayDecoder) + require.NoError(t, dec.SetData(len(values), data)) + + retained := make([]parquet.ByteArray, initial) + decoded, err := dec.Decode(retained) + require.NoError(t, err) + require.Equal(t, initial, decoded) + + discarded, err := dec.Discard(skip) + require.NoError(t, err) + require.Equal(t, min(skip, len(values)-initial), discarded) + next := initial + discarded + require.Equal(t, len(values)-next, dec.nvals) + require.Len(t, dec.lengths, len(values)-next) + require.Len(t, dec.prefixLengths, len(values)-next) + + rest := make([]parquet.ByteArray, len(values)+1) + decoded, err = dec.Decode(rest) + require.NoError(t, err) + require.Equal(t, len(values)-next, decoded) + for i, value := range rest[:decoded] { + require.Equal(t, values[next+i], string(value)) + } + for i, value := range retained { + require.Equal(t, values[i], string(value)) + } + + discarded, err = dec.Discard(1) + require.NoError(t, err) + require.Zero(t, discarded) + require.Zero(t, dec.nvals) + require.Empty(t, dec.lengths) + require.Empty(t, dec.prefixLengths) + require.Empty(t, dec.data) + }) + } + } +}