This is an automated email from the ASF dual-hosted git repository.
zeroshade pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-go.git
The following commit(s) were added to refs/heads/main by this push:
new 68c593e8 perf(parquet): reuse DELTA_BYTE_ARRAY discard storage (#1321)
68c593e8 is described below
commit 68c593e81b6d9b5d7c965966cd51405639dbab93
Author: Minh Vu <[email protected]>
AuthorDate: Fri Sep 25 20:27:15 2026 +0200
perf(parquet): reuse DELTA_BYTE_ARRAY discard storage (#1321)
**What**
- Reuse decoder-owned storage while `DELTA_BYTE_ARRAY` values are
discarded.
- Add correctness tests and discard benchmarks.
**Why**
- `Discard` currently allocates a new byte slice for every discarded
value with a non-empty suffix.
- Skipping 65,536 values creates 65,535 allocations only to throw the
values away.
**Implementation**
- Copy the first discarded value into owned scratch because the decoded
suffix can alias page data.
- Grow the scratch buffer amortized and reuse it across pages.
- Keep empty suffixes on the existing zero-copy path.
- Benchmark prefix-heavy and low-prefix data at 1,024 and 65,536 values.
The benchmark changes from 65,535 allocations to 0 allocations per
operation after warm-up, with about 3.3x lower discard time for 65,536
values.
Tests:
- `go test ./parquet/internal/encoding -count=1`
- `go test ./parquet/file -run
'^(TestWithEOFReader|TestInvalidHeaders|TestInvalidFooter|TestIncompleteMetadata|TestDeltaLengthByteArrayPackingWithNulls|TestDeltaBinaryPackedMultipleBatches|TestPageStreaming.*|TestPrimitiveReader|TestFullSeekRow|TestSkipEmptyRepeatedRows)$'
-count=1`
- `go test -race ./parquet/internal/encoding -run
'TestDeltaByteArrayDecoder(DiscardsAllEmptyValues|DiscardCopiesFirstValue|ReusesDiscardScratch|RejectsInvalidPrefixes|KeepsPartialDecodeResults)$'
-count=1`
- `go vet ./parquet/internal/encoding`
---
parquet/internal/encoding/delta_byte_array.go | 55 +++++++++----
.../delta_byte_array_decode_benchmark_test.go | 53 +++++++++++++
.../encoding/delta_byte_array_decode_test.go | 63 +++++++++++++++
.../delta_byte_array_discard_chunks_test.go | 72 +++++++++++++++++
.../delta_byte_array_discard_lifetime_test.go | 92 ++++++++++++++++++++++
.../internal/encoding/delta_length_byte_array.go | 15 +++-
6 files changed, 330 insertions(+), 20 deletions(-)
diff --git a/parquet/internal/encoding/delta_byte_array.go
b/parquet/internal/encoding/delta_byte_array.go
index 13428aca..3452b505 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,9 +163,10 @@ func (enc *DeltaByteArrayEncoder) FlushValues() (Buffer,
error) {
type DeltaByteArrayDecoder struct {
*DeltaLengthByteArrayDecoder
- prefixLengths []int32
- prefixScratch []int32
- lastVal parquet.ByteArray
+ prefixLengths []int32
+ prefixScratch []int32
+ lastVal parquet.ByteArray
+ discardScratch []byte
}
// Type returns the underlying physical type this decoder operates on, in this
case ByteArrays only
@@ -174,6 +176,29 @@ 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)
+ } 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)
+
+ // Discard roots reconstructed values at the start of discardScratch.
Its
+ // empty-suffix path only narrows lastVal, preserving that base pointer.
+ // Decode relies on this to detect aliases before later scratch reuse.
+ 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 {
@@ -229,15 +254,12 @@ func (d *DeltaByteArrayDecoder) Discard(n int) (int,
error) {
}
remaining := n
- tmp := make([]parquet.ByteArray, 1)
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.lastVal = tmp[0]
+ suffix := d.decodeOne()
+ d.setDiscardLastValue(nil, suffix)
d.prefixLengths = d.prefixLengths[1:]
remaining--
}
@@ -253,16 +275,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.decodeOne()
+ if len(suffix) == 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, suffix)
}
remaining--
}
@@ -370,6 +387,12 @@ func (d *DeltaByteArrayDecoder) Decode(out
[]parquet.ByteArray) (int, error) {
prefix := d.lastVal[:prefixLen:prefixLen]
if len(out[0]) == 0 {
+ // Discard can leave lastVal backed by reusable
discardScratch. Returning
+ // that prefix would let a later Discard overwrite data
the caller holds.
+ // setDiscardLastValue keeps lastVal at
discardScratch's base pointer.
+ 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_decode_benchmark_test.go
b/parquet/internal/encoding/delta_byte_array_decode_benchmark_test.go
index b03df3a0..0ff9a936 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 588e6b41..c839ee4e 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)
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 00000000..718a16ed
--- /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)
+ })
+ }
+ }
+}
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 00000000..f9f289de
--- /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]))
+}
diff --git a/parquet/internal/encoding/delta_length_byte_array.go
b/parquet/internal/encoding/delta_length_byte_array.go
index c06a3bb3..c9daf35f 100644
--- a/parquet/internal/encoding/delta_length_byte_array.go
+++ b/parquet/internal/encoding/delta_length_byte_array.go
@@ -165,6 +165,16 @@ func (d *DeltaLengthByteArrayDecoder) SetData(nvalues int,
data []byte) error {
return d.decoder.SetData(len(d.lengths), payload)
}
+// decodeOne decodes one value. The caller must ensure d.nvals > 0.
+func (d *DeltaLengthByteArrayDecoder) decodeOne() parquet.ByteArray {
+ length := d.lengths[0]
+ value := d.data[:length:length]
+ d.data = d.data[length:]
+ d.nvals--
+ d.lengths = d.lengths[1:]
+ return value
+}
+
func (d *DeltaLengthByteArrayDecoder) Discard(n int) (int, error) {
n = min(n, d.nvals)
for i := 0; i < n; i++ {
@@ -180,11 +190,8 @@ func (d *DeltaLengthByteArrayDecoder) Discard(n int) (int,
error) {
func (d *DeltaLengthByteArrayDecoder) Decode(out []parquet.ByteArray) (int,
error) {
max := utils.Min(len(out), d.nvals)
for i := 0; i < max; i++ {
- out[i] = d.data[:d.lengths[i]:d.lengths[i]]
- d.data = d.data[d.lengths[i]:]
+ out[i] = d.decodeOne()
}
- d.nvals -= max
- d.lengths = d.lengths[max:]
return max, nil
}