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 48f6b8a3 perf(parquet): decode DELTA_BYTE_ARRAY into contiguous 
storage (#1252)
48f6b8a3 is described below

commit 48f6b8a38ff1535a4c3a81ae21430085fff247cd
Author: Minh Vu <[email protected]>
AuthorDate: Mon Aug 31 20:26:17 2026 +0200

    perf(parquet): decode DELTA_BYTE_ARRAY into contiguous storage (#1252)
    
    ## Summary
    
    - Decode `DELTA_BYTE_ARRAY` suffixes in one batch.
    - Store reconstructed values in one contiguous byte arena per `Decode`
    call.
    - Keep prefix validation before allocation and reconstruction.
    - Keep separate arenas for separate decode calls so earlier results stay
    valid.
    - Add coverage for partial reads, spaced reads, empty values, page
    resets, and invalid state.
    - Add benchmarks for high and low prefix sharing with 1K and 64K values.
    
    ## Benchmark
    
    Apple M1 Pro. Compared with main at `6b039a76`. The benchmark uses
    65,536 values, 1 second per sample, 5 samples, and `GOMAXPROCS=1`.
    
    | case | main | branch | change | allocs/op |
    | --- | ---: | ---: | ---: | ---: |
    | prefix-heavy, batch 128 | 3.678 ms | 2.714 ms | 26.2% faster | 65,553
    -> 529 |
    | prefix-heavy, full page | 3.625 ms | 2.727 ms | 24.8% faster | 65,551
    -> 17 |
    | low-prefix, batch 128 | 3.172 ms | 2.424 ms | 23.6% faster | 65,552 ->
    529 |
    | low-prefix, full page | 3.169 ms | 2.382 ms | 24.8% faster | 65,551 ->
    17 |
    
    Command:
    
    ```text
    GOMAXPROCS=1 go test ./parquet/internal/encoding -run '^$' -bench 
'^BenchmarkDeltaByteArrayDecoderDecode/(prefix-heavy|low-prefix)/65536/batch-(128|65536)$'
 -benchmem -benchtime=1s -count=5 -cpu=1
    ```
    
    ## Checks
    
    - `go test ./parquet/internal/encoding -count=1`
    - `go test -race ./parquet/internal/encoding -count=1`
    - `go vet ./parquet/internal/encoding`
    - `go test ./parquet/... -run '^$'`
    - `git diff --check`
---
 parquet/internal/encoding/delta_byte_array.go      | 119 +++++++++---
 .../delta_byte_array_decode_benchmark_test.go      | 109 +++++++++++
 .../encoding/delta_byte_array_decode_test.go       | 214 +++++++++++++++++++++
 3 files changed, 415 insertions(+), 27 deletions(-)

diff --git a/parquet/internal/encoding/delta_byte_array.go 
b/parquet/internal/encoding/delta_byte_array.go
index 86b7d585..e6c641b3 100644
--- a/parquet/internal/encoding/delta_byte_array.go
+++ b/parquet/internal/encoding/delta_byte_array.go
@@ -257,6 +257,64 @@ func (d *DeltaByteArrayDecoder) Discard(n int) (int, 
error) {
        return n, nil
 }
 
+func (d *DeltaByteArrayDecoder) decodedArenaSize(max int) (int, error) {
+       maxInt := int(^uint(0) >> 1)
+       total := 0
+       prefixLengths := d.prefixLengths
+       suffixLengths := d.lengths
+       previousLen := len(d.lastVal)
+       if d.lastVal == nil {
+               if len(prefixLengths) == 0 || prefixLengths[0] != 0 {
+                       return 0, errors.New("parquet: first delta byte array 
prefix length must be zero")
+               }
+               if len(suffixLengths) == 0 {
+                       return 0, errors.New("parquet: not enough delta byte 
array suffix lengths")
+               }
+
+               suffixLen := suffixLengths[0]
+               if suffixLen < 0 {
+                       return 0, fmt.Errorf("parquet: negative delta byte 
array length %d", suffixLen)
+               }
+               previousLen = int(suffixLen)
+               total = previousLen
+               prefixLengths = prefixLengths[1:]
+               suffixLengths = suffixLengths[1:]
+               max--
+       }
+
+       for i := 0; i < max; i++ {
+               if i >= len(prefixLengths) {
+                       return 0, errors.New("parquet: not enough delta byte 
array prefix lengths")
+               }
+               if i >= len(suffixLengths) {
+                       return 0, errors.New("parquet: not enough delta byte 
array suffix lengths")
+               }
+               prefixLen := prefixLengths[i]
+               suffixLen := suffixLengths[i]
+               if prefixLen < 0 || int(prefixLen) > previousLen {
+                       return 0, fmt.Errorf("parquet: invalid delta byte array 
prefix length %d", prefixLen)
+               }
+               if suffixLen < 0 {
+                       return 0, fmt.Errorf("parquet: negative delta byte 
array length %d", suffixLen)
+               }
+
+               valueLen := int(prefixLen)
+               if int(suffixLen) > maxInt-valueLen {
+                       return 0, errors.New("parquet: delta byte array value 
length overflows int")
+               }
+               valueLen += int(suffixLen)
+               if suffixLen != 0 {
+                       if valueLen > maxInt-total {
+                               return 0, errors.New("parquet: decoded delta 
byte array size overflows int")
+                       }
+                       total += valueLen
+               }
+               previousLen = valueLen
+       }
+
+       return total, nil
+}
+
 // Decode decodes byte arrays into the slice provided and returns the number 
of values actually decoded
 func (d *DeltaByteArrayDecoder) Decode(out []parquet.ByteArray) (int, error) {
        max := utils.Min(len(out), d.nvals)
@@ -265,45 +323,52 @@ func (d *DeltaByteArrayDecoder) Decode(out 
[]parquet.ByteArray) (int, error) {
        }
        out = out[:max]
 
-       var err error
+       arenaSize, err := d.decodedArenaSize(max)
+       if err != nil {
+               return 0, err
+       }
+       arena := make([]byte, arenaSize)
+       arenaOffset := 0
+       decoded, err := d.DeltaLengthByteArrayDecoder.Decode(out)
+       if err != nil {
+               return 0, err
+       }
+       if decoded != max {
+               return 0, errors.New("parquet: not enough delta byte array 
suffix values")
+       }
+
        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")
-               }
-               _, err = d.DeltaLengthByteArrayDecoder.Decode(out[:1])
-               if err != nil {
-                       return 0, err
-               }
-               d.lastVal = out[0]
-               out = out[1:]
+               valueLen := len(out[0])
+               value := arena[arenaOffset : arenaOffset+valueLen : 
arenaOffset+valueLen]
+               copy(value, out[0])
+               out[0] = value
+               d.lastVal = value
+               arenaOffset += valueLen
                d.prefixLengths = d.prefixLengths[1:]
+               out = out[1:]
        }
 
-       var prefixLen int32
-       suffixHolder := make([]parquet.ByteArray, 1)
        for len(out) > 0 {
                if len(d.prefixLengths) == 0 {
                        return 0, errors.New("parquet: not enough delta byte 
array prefix lengths")
                }
-               prefixLen, d.prefixLengths = d.prefixLengths[0], 
d.prefixLengths[1:]
-               if prefixLen < 0 || int(prefixLen) > len(d.lastVal) {
-                       return 0, fmt.Errorf("parquet: invalid delta byte array 
prefix length %d", prefixLen)
-               }
+               prefixLen := d.prefixLengths[0]
+               d.prefixLengths = d.prefixLengths[1:]
 
                prefix := d.lastVal[:prefixLen:prefixLen]
-               _, err = d.DeltaLengthByteArrayDecoder.Decode(suffixHolder)
-               if err != nil {
-                       return 0, err
-               }
-
-               if len(suffixHolder[0]) == 0 {
+               if len(out[0]) == 0 {
                        d.lastVal = prefix
-               } else {
-                       d.lastVal = make([]byte, 
int(prefixLen)+len(suffixHolder[0]))
-                       copy(d.lastVal, prefix)
-                       copy(d.lastVal[prefixLen:], suffixHolder[0])
+                       out[0], out = prefix, out[1:]
+                       continue
                }
-               out[0], out = d.lastVal, out[1:]
+
+               valueLen := int(prefixLen) + len(out[0])
+               value := arena[arenaOffset : arenaOffset+valueLen : 
arenaOffset+valueLen]
+               copy(value, prefix)
+               copy(value[len(prefix):], out[0])
+               out[0], out = value, out[1:]
+               d.lastVal = value
+               arenaOffset += valueLen
        }
        return max, nil
 }
diff --git 
a/parquet/internal/encoding/delta_byte_array_decode_benchmark_test.go 
b/parquet/internal/encoding/delta_byte_array_decode_benchmark_test.go
new file mode 100644
index 00000000..b03df3a0
--- /dev/null
+++ b/parquet/internal/encoding/delta_byte_array_decode_benchmark_test.go
@@ -0,0 +1,109 @@
+// 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"
+       "strings"
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/parquet"
+)
+
+func BenchmarkDeltaByteArrayDecoderDecode(b *testing.B) {
+       repeatedValue := strings.Repeat("x", 1024)
+       for _, test := range []struct {
+               name  string
+               value func(int) string
+       }{
+               {
+                       name: "identical",
+                       value: func(int) string {
+                               return repeatedValue
+                       },
+               },
+               {
+                       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)
+
+                               for _, batchSize := range []int{128, 1024, 
nvalues} {
+                                       batchSize := batchSize
+                                       b.Run(fmt.Sprintf("batch-%d", 
batchSize), func(b *testing.B) {
+                                               output := 
make([]parquet.ByteArray, batchSize)
+                                               dec := 
NewDecoder(parquet.Types.ByteArray, parquet.Encodings.DeltaByteArray,
+                                                       nil, 
memory.DefaultAllocator).(ByteArrayDecoder)
+                                               b.SetBytes(int64(inputBytes))
+                                               b.ReportAllocs()
+                                               b.ResetTimer()
+                                               for b.Loop() {
+                                                       if err := 
dec.SetData(nvalues, encoded); err != nil {
+                                                               b.Fatal(err)
+                                                       }
+                                                       remaining := nvalues
+                                                       for remaining > 0 {
+                                                               count := 
min(batchSize, remaining)
+                                                               decoded, err := 
dec.Decode(output[:count])
+                                                               if err != nil {
+                                                                       
b.Fatal(err)
+                                                               }
+                                                               if decoded != 
count {
+                                                                       
b.Fatalf("decoded %d values, expected %d", decoded, count)
+                                                               }
+                                                               remaining -= 
count
+                                                       }
+                                               }
+                                       })
+                               }
+                       })
+               }
+       }
+}
+
+func encodeDeltaByteArrayValues(values []parquet.ByteArray) []byte {
+       enc := NewEncoder(parquet.Types.ByteArray, 
parquet.Encodings.DeltaByteArray,
+               false, nil, memory.DefaultAllocator).(ByteArrayEncoder)
+       defer enc.Release()
+       enc.Put(values)
+       buf, err := enc.FlushValues()
+       if err != nil {
+               panic(err)
+       }
+       defer buf.Release()
+       return append([]byte(nil), buf.Bytes()...)
+}
diff --git a/parquet/internal/encoding/delta_byte_array_decode_test.go 
b/parquet/internal/encoding/delta_byte_array_decode_test.go
new file mode 100644
index 00000000..1a5ea5f2
--- /dev/null
+++ b/parquet/internal/encoding/delta_byte_array_decode_test.go
@@ -0,0 +1,214 @@
+// 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"
+       "strings"
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/parquet"
+       "github.com/stretchr/testify/require"
+)
+
+func encodeDeltaByteArrayPage(t *testing.T, values []string) []byte {
+       t.Helper()
+
+       input := make([]parquet.ByteArray, len(values))
+       for i, value := range values {
+               input[i] = parquet.ByteArray(value)
+       }
+
+       enc := NewEncoder(parquet.Types.ByteArray, 
parquet.Encodings.DeltaByteArray,
+               false, nil, memory.DefaultAllocator).(ByteArrayEncoder)
+       defer enc.Release()
+       enc.Put(input)
+
+       buf, err := enc.FlushValues()
+       require.NoError(t, err)
+       defer buf.Release()
+       return append([]byte(nil), buf.Bytes()...)
+}
+
+func requireDecodedStrings(t *testing.T, got []parquet.ByteArray, want 
[]string) {
+       t.Helper()
+       require.Len(t, got, len(want))
+       for i, value := range want {
+               require.Equal(t, value, string(got[i]), "value %d", i)
+       }
+}
+
+func TestDeltaByteArrayDecoderKeepsPartialDecodeResults(t *testing.T) {
+       values := []string{
+               "", "partition/000/value/000", "partition/000/value/001",
+               "partition/001/value/000", "partition/001/value/001", 
"partition/001/value/001",
+       }
+       data := encodeDeltaByteArrayPage(t, values)
+       dec := NewDecoder(parquet.Types.ByteArray, 
parquet.Encodings.DeltaByteArray,
+               nil, memory.DefaultAllocator).(ByteArrayDecoder)
+       require.NoError(t, dec.SetData(len(values), data))
+
+       decoded := make([]parquet.ByteArray, len(values))
+       for i := range decoded {
+               n, err := dec.Decode(decoded[i : i+1])
+               require.NoError(t, err)
+               require.Equal(t, 1, n)
+       }
+
+       requireDecodedStrings(t, decoded, values)
+}
+
+func TestDeltaByteArrayDecoderDecodesAllEmptyValues(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))
+
+       decoded := make([]parquet.ByteArray, len(values))
+       for i := range decoded {
+               n, err := dec.Decode(decoded[i : i+1])
+               require.NoError(t, err)
+               require.Equal(t, 1, n)
+       }
+
+       for i, value := range decoded {
+               require.Empty(t, value, "value %d", i)
+       }
+}
+
+func TestDeltaByteArrayDecoderReusesValuesWithoutSuffixes(t *testing.T) {
+       value := strings.Repeat("x", 64*1024)
+       values := make([]string, 128)
+       for i := range values {
+               values[i] = value
+       }
+       values = append(values, value[:1024], "", "new value", "new value")
+       data := encodeDeltaByteArrayPage(t, values)
+
+       for _, batchSize := range []int{1, 17, len(values)} {
+               t.Run(fmt.Sprintf("batch-%d", batchSize), func(t *testing.T) {
+                       dec := NewDecoder(parquet.Types.ByteArray, 
parquet.Encodings.DeltaByteArray,
+                               nil, 
memory.DefaultAllocator).(*DeltaByteArrayDecoder)
+                       require.NoError(t, dec.SetData(len(values), data))
+
+                       arenaSize, err := dec.decodedArenaSize(len(values))
+                       require.NoError(t, err)
+                       require.Equal(t, len(value)+len("new value"), arenaSize)
+
+                       decoded := make([]parquet.ByteArray, len(values))
+                       for offset := 0; offset < len(values); offset += 
batchSize {
+                               end := min(offset+batchSize, len(values))
+                               n, err := dec.Decode(decoded[offset:end])
+                               require.NoError(t, err)
+                               require.Equal(t, end-offset, n)
+                       }
+                       requireDecodedStrings(t, decoded, values)
+                       for i := 1; i < 128; i++ {
+                               require.Same(t, &decoded[0][0], &decoded[i][0])
+                       }
+                       require.Equal(t, len(decoded[128]), cap(decoded[128]))
+               })
+       }
+}
+
+func TestDeltaByteArrayDecoderKeepsResultsAcrossPages(t *testing.T) {
+       firstValues := []string{"first/000", "first/001", "first/002"}
+       secondValues := []string{"second/000", "second/001"}
+       dec := NewDecoder(parquet.Types.ByteArray, 
parquet.Encodings.DeltaByteArray,
+               nil, memory.DefaultAllocator).(ByteArrayDecoder)
+
+       firstData := encodeDeltaByteArrayPage(t, firstValues)
+       require.NoError(t, dec.SetData(len(firstValues), firstData))
+       firstOut := make([]parquet.ByteArray, len(firstValues))
+       decoded, err := dec.Decode(firstOut)
+       require.NoError(t, err)
+       require.Equal(t, len(firstValues), decoded)
+
+       secondData := encodeDeltaByteArrayPage(t, secondValues)
+       require.NoError(t, dec.SetData(len(secondValues), secondData))
+       secondOut := make([]parquet.ByteArray, len(secondValues))
+       decoded, err = dec.Decode(secondOut)
+       require.NoError(t, err)
+       require.Equal(t, len(secondValues), decoded)
+
+       requireDecodedStrings(t, firstOut, firstValues)
+       requireDecodedStrings(t, secondOut, secondValues)
+}
+
+func TestDeltaByteArrayDecoderDecodeSpaced(t *testing.T) {
+       values := []string{"a/000", "a/001", "b/000", "b/001"}
+       data := encodeDeltaByteArrayPage(t, values)
+       dec := NewDecoder(parquet.Types.ByteArray, 
parquet.Encodings.DeltaByteArray,
+               nil, memory.DefaultAllocator).(ByteArrayDecoder)
+       require.NoError(t, dec.SetData(len(values), data))
+
+       validBits := []byte{0b00101101}
+       out := make([]parquet.ByteArray, 6)
+       decoded, err := dec.DecodeSpaced(out, 2, validBits, 0)
+       require.NoError(t, err)
+       require.Equal(t, len(out), decoded)
+       for i, value := range map[int]string{0: "a/000", 2: "a/001", 3: 
"b/000", 5: "b/001"} {
+               require.Equal(t, value, string(out[i]), "value %d", i)
+       }
+}
+
+func TestDeltaByteArrayDecoderDecodedSizeRejectsInvalidState(t *testing.T) {
+       tests := []struct {
+               name          string
+               lastVal       parquet.ByteArray
+               prefixLengths []int32
+               lengths       []int32
+               want          string
+       }{
+               {
+                       name:          "nonzero first prefix",
+                       prefixLengths: []int32{1},
+                       lengths:       []int32{1},
+                       want:          "first delta byte array prefix length 
must be zero",
+               },
+               {
+                       name:          "prefix beyond previous value",
+                       lastVal:       parquet.ByteArray("abc"),
+                       prefixLengths: []int32{4},
+                       lengths:       []int32{1},
+                       want:          "invalid delta byte array prefix length 
4",
+               },
+               {
+                       name:          "negative suffix",
+                       lastVal:       parquet.ByteArray("abc"),
+                       prefixLengths: []int32{1},
+                       lengths:       []int32{-1},
+                       want:          "negative delta byte array length -1",
+               },
+       }
+
+       for _, tt := range tests {
+               t.Run(tt.name, func(t *testing.T) {
+                       dec := &DeltaByteArrayDecoder{
+                               DeltaLengthByteArrayDecoder: 
&DeltaLengthByteArrayDecoder{
+                                       lengths: tt.lengths,
+                               },
+                               prefixLengths: tt.prefixLengths,
+                               lastVal:       tt.lastVal,
+                       }
+                       _, err := dec.decodedArenaSize(1)
+                       require.EqualError(t, err, fmt.Sprintf("parquet: %s", 
tt.want))
+               })
+       }
+}

Reply via email to