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 3d7bda19 perf(parquet): decode BSS FLBA into contiguous storage (#1172)
3d7bda19 is described below

commit 3d7bda19bc085ba66caa98414e4ca1cc1feee8f0
Author: Minh Vu <[email protected]>
AuthorDate: Fri Aug 28 23:00:13 2026 +0200

    perf(parquet): decode BSS FLBA into contiguous storage (#1172)
    
    ### Rationale for this change
    
    `BYTE_STREAM_SPLIT` decoding for `FIXED_LEN_BYTE_ARRAY` currently
    allocates a separate byte slice for every output value when the output
    has no reusable capacity. A 65,536-value batch therefore performs 65,536
    decoder allocations.
    
    ### What changes are included in this PR?
    
    - Allocate one contiguous backing block for output values that need
    storage.
    - Keep reusing output slices that already have enough capacity.
    - Give each decoded value an exact-capacity slice.
    - Only prepare the output prefix that is actually decoded.
    - Add coverage for contiguous storage, mixed reusable output, output
    lifetime across decoder resets, partial output, and spaced decoding.
    - Add cold-output benchmarks for widths 2, 4, 8, 16, and 32.
    
    Apple M1 Pro, GOMAXPROCS=1, 65,536 values:
    
    | Width | Before | After | Change | Allocations |
    | ---: | ---: | ---: | ---: | ---: |
    | 16 | 2.42 ms | 1.55 ms | -35.7% | 65,537 to 2 |
    | 32 | 3.82 ms | 2.61 ms | -31.7% | 65,537 to 2 |
    
    Bytes allocated stay unchanged. The existing reusable-output benchmark
    also remains allocation-free and improved by about 8%.
    
    ### Are these changes tested?
    
    - `PARQUET_TEST_DATA=$PWD/parquet-testing/data go test ./parquet/...`
    - `go test -race ./parquet/internal/encoding`
    - `go vet -composites=false ./parquet/internal/encoding`
    - Cross-compiled the encoding package tests for linux/amd64 and
    linux/s390x.
    
    ### Are there any user-facing changes?
    
    No.
---
 .../internal/encoding/encoding_benchmarks_test.go  |  44 +++++
 .../encoding/fixed_len_byte_array_decoder.go       |  32 +++-
 .../encoding/fixed_len_byte_array_decoder_test.go  | 186 +++++++++++++++++++++
 3 files changed, 257 insertions(+), 5 deletions(-)

diff --git a/parquet/internal/encoding/encoding_benchmarks_test.go 
b/parquet/internal/encoding/encoding_benchmarks_test.go
index 781967a9..7946e844 100644
--- a/parquet/internal/encoding/encoding_benchmarks_test.go
+++ b/parquet/internal/encoding/encoding_benchmarks_test.go
@@ -714,6 +714,50 @@ func BenchmarkByteStreamSplitDecodingFixedLenByteArray(b 
*testing.B) {
        }
 }
 
+func BenchmarkByteStreamSplitDecodingFixedLenByteArrayColdOutput(b *testing.B) 
{
+       for _, width := range []int{2, 4, 8, 16, 32} {
+               for _, size := range []int{MINSIZE, MAXSIZE} {
+                       b.Run(fmt.Sprintf("width %d/len %d", width, size), 
func(b *testing.B) {
+                               values := make([]parquet.FixedLenByteArray, 
size)
+                               for idx := range values {
+                                       values[idx] = 
make(parquet.FixedLenByteArray, width)
+                                       for byteIdx := range values[idx] {
+                                               values[idx][byteIdx] = byte(idx 
+ byteIdx)
+                                       }
+                               }
+
+                               col := 
schema.NewColumn(schema.NewFixedLenByteArrayNode("fixedlenbytearray", 
parquet.Repetitions.Required, int32(width), -1), 0, 0)
+                               encoder := 
encoding.NewEncoder(parquet.Types.FixedLenByteArray, 
parquet.Encodings.ByteStreamSplit,
+                                       false, col, 
memory.DefaultAllocator).(encoding.FixedLenByteArrayEncoder)
+                               encoder.Put(values)
+                               buf, err := encoder.FlushValues()
+                               if err != nil {
+                                       b.Fatal(err)
+                               }
+                               defer buf.Release()
+
+                               decoder := 
encoding.NewDecoder(parquet.Types.FixedLenByteArray, 
parquet.Encodings.ByteStreamSplit, col, 
memory.DefaultAllocator).(encoding.FixedLenByteArrayDecoder)
+                               b.ReportAllocs()
+                               b.SetBytes(int64(size * width))
+                               b.ResetTimer()
+                               for b.Loop() {
+                                       output := 
make([]parquet.FixedLenByteArray, size)
+                                       if err := decoder.SetData(size, 
buf.Bytes()); err != nil {
+                                               b.Fatal(err)
+                                       }
+                                       decoded, err := decoder.Decode(output)
+                                       if err != nil {
+                                               b.Fatal(err)
+                                       }
+                                       if decoded != size {
+                                               b.Fatalf("decoded %d values, 
want %d", decoded, size)
+                                       }
+                               }
+                       })
+               }
+       }
+}
+
 func BenchmarkDeltaBinaryPackedEncodingInt32(b *testing.B) {
        for sz := MINSIZE; sz < MAXSIZE+1; sz *= 2 {
                b.Run(fmt.Sprintf("len %d", sz), func(b *testing.B) {
diff --git a/parquet/internal/encoding/fixed_len_byte_array_decoder.go 
b/parquet/internal/encoding/fixed_len_byte_array_decoder.go
index 89c47632..080e9a06 100644
--- a/parquet/internal/encoding/fixed_len_byte_array_decoder.go
+++ b/parquet/internal/encoding/fixed_len_byte_array_decoder.go
@@ -173,12 +173,13 @@ func (dec *ByteStreamSplitFixedLenByteArrayDecoder) 
Decode(out []parquet.FixedLe
                return 0, errors.New("parquet: eof exception")
        }
 
-       for i := range out {
-               if cap(out[i]) < dec.typeLen {
-                       out[i] = make(parquet.FixedLenByteArray, dec.typeLen)
-               } else {
-                       out[i] = out[i][:dec.typeLen]
+       out = out[:toRead]
+       for idx := range out {
+               if cap(out[idx]) < dec.typeLen {
+                       dec.prepareOutput(out[idx:])
+                       break
                }
+               out[idx] = out[idx][:dec.typeLen]
        }
 
        switch dec.typeLen {
@@ -197,6 +198,27 @@ func (dec *ByteStreamSplitFixedLenByteArrayDecoder) 
Decode(out []parquet.FixedLe
        return toRead, nil
 }
 
+// prepareOutput allocates storage for the entries in out that do not have 
enough
+// capacity, while continuing to reuse the entries that do.
+func (dec *ByteStreamSplitFixedLenByteArrayDecoder) prepareOutput(out 
[]parquet.FixedLenByteArray) {
+       missing := 0
+       for idx := range out {
+               if cap(out[idx]) < dec.typeLen {
+                       missing++
+               }
+       }
+
+       storage := make([]byte, missing*dec.typeLen)
+       for idx := range out {
+               if cap(out[idx]) < dec.typeLen {
+                       out[idx] = storage[:dec.typeLen:dec.typeLen]
+                       storage = storage[dec.typeLen:]
+               } else {
+                       out[idx] = out[idx][:dec.typeLen]
+               }
+       }
+}
+
 func (dec *ByteStreamSplitFixedLenByteArrayDecoder) DecodeSpaced(out 
[]parquet.FixedLenByteArray, nullCount int, validBits []byte, validBitsOffset 
int64) (int, error) {
        toRead := len(out) - nullCount
        valuesRead, err := dec.Decode(out[:toRead])
diff --git a/parquet/internal/encoding/fixed_len_byte_array_decoder_test.go 
b/parquet/internal/encoding/fixed_len_byte_array_decoder_test.go
new file mode 100644
index 00000000..732d61e7
--- /dev/null
+++ b/parquet/internal/encoding/fixed_len_byte_array_decoder_test.go
@@ -0,0 +1,186 @@
+// 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"
+       "unsafe"
+
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/parquet"
+       "github.com/apache/arrow-go/v18/parquet/schema"
+       "github.com/stretchr/testify/require"
+)
+
+func TestByteStreamSplitFixedLenByteArrayDecoderContiguousOutput(t *testing.T) 
{
+       for _, width := range []int{2, 4, 8, 16, 32} {
+               t.Run(fmt.Sprintf("width=%d", width), func(t *testing.T) {
+                       values := makeFixedLenByteArrayValues(8, width, 0)
+                       decoder := 
newByteStreamSplitFixedLenByteArrayDecoder(t, width, values)
+
+                       out := make([]parquet.FixedLenByteArray, len(values))
+                       decoded, err := decoder.Decode(out)
+                       require.NoError(t, err)
+                       require.Equal(t, len(values), decoded)
+                       require.Equal(t, values, out)
+
+                       for idx := range out {
+                               require.Equal(t, width, cap(out[idx]))
+                               if idx > 0 {
+                                       previous := 
uintptr(unsafe.Pointer(&out[idx-1][0]))
+                                       current := 
uintptr(unsafe.Pointer(&out[idx][0]))
+                                       require.Equal(t, uintptr(width), 
current-previous)
+                               }
+                       }
+               })
+       }
+}
+
+func TestByteStreamSplitFixedLenByteArrayDecoderReusesProvidedOutput(t 
*testing.T) {
+       const width = 16
+       values := makeFixedLenByteArrayValues(4, width, 0)
+       decoder := newByteStreamSplitFixedLenByteArrayDecoder(t, width, values)
+
+       out := []parquet.FixedLenByteArray{
+               make([]byte, 0, width+4),
+               nil,
+               make([]byte, width),
+               nil,
+       }
+       firstPtr := unsafe.Pointer(unsafe.SliceData(out[0]))
+       thirdPtr := unsafe.Pointer(unsafe.SliceData(out[2]))
+
+       decoded, err := decoder.Decode(out)
+       require.NoError(t, err)
+       require.Equal(t, len(values), decoded)
+       require.Equal(t, values, out)
+       require.Equal(t, firstPtr, unsafe.Pointer(unsafe.SliceData(out[0])))
+       require.Equal(t, thirdPtr, unsafe.Pointer(unsafe.SliceData(out[2])))
+       require.Equal(t, uintptr(width),
+               
uintptr(unsafe.Pointer(&out[3][0]))-uintptr(unsafe.Pointer(&out[1][0])))
+}
+
+func TestByteStreamSplitFixedLenByteArrayDecoderMixedOutputAllocations(t 
*testing.T) {
+       const width = 16
+       values := makeFixedLenByteArrayValues(4, width, 0)
+       data := encodeByteStreamSplitFixedLenByteArray(values, width)
+       decoder := newByteStreamSplitFixedLenByteArrayDecoder(t, width, values)
+
+       out := []parquet.FixedLenByteArray{
+               make([]byte, width),
+               nil,
+               make([]byte, width),
+               nil,
+       }
+       allocs := testing.AllocsPerRun(100, func() {
+               out[1] = nil
+               out[3] = nil
+               require.NoError(t, decoder.SetData(len(values), data))
+               _, err := decoder.Decode(out)
+               require.NoError(t, err)
+       })
+       require.Equal(t, float64(1), allocs)
+}
+
+func TestByteStreamSplitFixedLenByteArrayDecoderKeepsPreviousOutput(t 
*testing.T) {
+       const width = 16
+       firstValues := makeFixedLenByteArrayValues(4, width, 0)
+       decoder := newByteStreamSplitFixedLenByteArrayDecoder(t, width, 
firstValues)
+
+       firstOut := make([]parquet.FixedLenByteArray, len(firstValues))
+       _, err := decoder.Decode(firstOut)
+       require.NoError(t, err)
+
+       secondValues := makeFixedLenByteArrayValues(4, width, 100)
+       require.NoError(t, decoder.SetData(len(secondValues), 
encodeByteStreamSplitFixedLenByteArray(secondValues, width)))
+       secondOut := make([]parquet.FixedLenByteArray, len(secondValues))
+       _, err = decoder.Decode(secondOut)
+       require.NoError(t, err)
+
+       require.Equal(t, firstValues, firstOut)
+       require.Equal(t, secondValues, secondOut)
+}
+
+func TestByteStreamSplitFixedLenByteArrayDecoderPartialOutput(t *testing.T) {
+       const width = 8
+       values := makeFixedLenByteArrayValues(3, width, 0)
+       decoder := newByteStreamSplitFixedLenByteArrayDecoder(t, width, values)
+
+       sentinel := parquet.FixedLenByteArray("sentinel")
+       out := make([]parquet.FixedLenByteArray, 5)
+       out[3] = sentinel
+       out[4] = sentinel
+
+       decoded, err := decoder.Decode(out)
+       require.NoError(t, err)
+       require.Equal(t, len(values), decoded)
+       require.Equal(t, values, out[:decoded])
+       require.Equal(t, unsafe.Pointer(&sentinel[0]), 
unsafe.Pointer(&out[3][0]))
+       require.Equal(t, unsafe.Pointer(&sentinel[0]), 
unsafe.Pointer(&out[4][0]))
+}
+
+func TestByteStreamSplitFixedLenByteArrayDecoderSpacedOutput(t *testing.T) {
+       const width = 8
+       values := makeFixedLenByteArrayValues(3, width, 0)
+       decoder := newByteStreamSplitFixedLenByteArrayDecoder(t, width, values)
+       out := make([]parquet.FixedLenByteArray, 5)
+
+       decoded, err := decoder.DecodeSpaced(out, 2, []byte{0b00010101}, 0)
+       require.NoError(t, err)
+       require.Equal(t, len(out), decoded)
+       require.Equal(t, values[0], out[0])
+       require.Equal(t, values[1], out[2])
+       require.Equal(t, values[2], out[4])
+
+       first := uintptr(unsafe.Pointer(&out[0][0]))
+       second := uintptr(unsafe.Pointer(&out[2][0]))
+       third := uintptr(unsafe.Pointer(&out[4][0]))
+       require.Equal(t, uintptr(width), second-first)
+       require.Equal(t, uintptr(width), third-second)
+}
+
+func newByteStreamSplitFixedLenByteArrayDecoder(t *testing.T, width int, 
values []parquet.FixedLenByteArray) FixedLenByteArrayDecoder {
+       t.Helper()
+
+       node := schema.NewFixedLenByteArrayNode("value", 
parquet.Repetitions.Required, int32(width), -1)
+       column := schema.NewColumn(node, 0, 0)
+       decoder := NewDecoder(parquet.Types.FixedLenByteArray, 
parquet.Encodings.ByteStreamSplit, column, 
memory.DefaultAllocator).(FixedLenByteArrayDecoder)
+       require.NoError(t, decoder.SetData(len(values), 
encodeByteStreamSplitFixedLenByteArray(values, width)))
+       return decoder
+}
+
+func makeFixedLenByteArrayValues(length, width int, offset byte) 
[]parquet.FixedLenByteArray {
+       values := make([]parquet.FixedLenByteArray, length)
+       for valueIdx := range values {
+               values[valueIdx] = make(parquet.FixedLenByteArray, width)
+               for byteIdx := range values[valueIdx] {
+                       values[valueIdx][byteIdx] = offset + 
byte(valueIdx*width+byteIdx)
+               }
+       }
+       return values
+}
+
+func encodeByteStreamSplitFixedLenByteArray(values 
[]parquet.FixedLenByteArray, width int) []byte {
+       data := make([]byte, len(values)*width)
+       for valueIdx, value := range values {
+               for byteIdx, valueByte := range value {
+                       data[byteIdx*len(values)+valueIdx] = valueByte
+               }
+       }
+       return data
+}

Reply via email to