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 7efe1c01 perf(parquet): reuse DELTA byte-array scratch (#1261)
7efe1c01 is described below

commit 7efe1c0125eeaedd4a3da8b5b2bc7d5b0c6b96fd
Author: Minh Vu <[email protected]>
AuthorDate: Mon Aug 31 22:39:21 2026 +0200

    perf(parquet): reuse DELTA byte-array scratch (#1261)
    
    ## Summary
    
    - Reuse `[]parquet.ByteArray` scratch for spaced compaction in
    `DELTA_LENGTH_BYTE_ARRAY`.
    - Reuse the same scratch pattern in `DELTA_BYTE_ARRAY`.
    - Clear the scratch after encoding so input payload slices are not
    retained.
    - Add offset, correctness, reuse, and benchmark coverage.
    
    ## Benchmark
    
    64K values with 16-byte values on an Apple M1 Pro, Go 1.26.3, 1 CPU, 250
    ms per sample, 9 samples. Values are medians.
    
    | Encoder / validity | ns/op before | ns/op after | B/op before | B/op
    after | allocs/op before | allocs/op after |
    | --- | ---: | ---: | ---: | ---: | ---: | ---: |
    | DELTA_LENGTH_BYTE_ARRAY / all valid | 1,209,684 | 886,099 | 1,577,916
    | 5,000 | 2,055 | 2,054 |
    | DELTA_LENGTH_BYTE_ARRAY / 50% null | 1,049,068 | 692,072 | 1,575,636 |
    2,696 | 1,031 | 1,030 |
    | DELTA_BYTE_ARRAY / all valid | 2,618,095 | 1,957,687 | 1,633,022 |
    57,918 | 4,111 | 4,109 |
    | DELTA_BYTE_ARRAY / 50% null | 1,776,488 | 1,141,407 | 1,604,656 |
    31,066 | 2,063 | 2,061 |
    
    ## Checks
    
    - `go test ./...`
    - `go test -race ./parquet/internal/encoding -count=1`
    - `go vet ./parquet/internal/encoding`
    - `git diff --check`
    
    No public API changes.
---
 parquet/internal/encoding/delta_byte_array.go      |  12 +-
 .../delta_byte_array_spaced_benchmark_test.go      |  84 +++++++++++++
 .../encoding/delta_byte_array_spaced_test.go       | 135 +++++++++++++++++++++
 .../internal/encoding/delta_length_byte_array.go   |  12 +-
 4 files changed, 237 insertions(+), 6 deletions(-)

diff --git a/parquet/internal/encoding/delta_byte_array.go 
b/parquet/internal/encoding/delta_byte_array.go
index e6c641b3..b2ed5bf5 100644
--- a/parquet/internal/encoding/delta_byte_array.go
+++ b/parquet/internal/encoding/delta_byte_array.go
@@ -40,6 +40,7 @@ type DeltaByteArrayEncoder struct {
 
        prefixLengths [deltaByteArrayBatchSize]int32
        suffixes      [deltaByteArrayBatchSize]parquet.ByteArray
+       spacedScratch []parquet.ByteArray
 
        lastVal parquet.ByteArray
 }
@@ -118,9 +119,14 @@ func (enc *DeltaByteArrayEncoder) Put(in 
[]parquet.ByteArray) {
 // to compress the data before writing it without the null slots.
 func (enc *DeltaByteArrayEncoder) PutSpaced(in []parquet.ByteArray, validBits 
[]byte, validBitsOffset int64) {
        if validBits != nil {
-               data := make([]parquet.ByteArray, len(in))
-               nvalid := spacedCompress(in, data, validBits, validBitsOffset)
-               enc.Put(data[:nvalid])
+               if cap(enc.spacedScratch) < len(in) {
+                       enc.spacedScratch = make([]parquet.ByteArray, len(in))
+               } else {
+                       enc.spacedScratch = enc.spacedScratch[:len(in)]
+               }
+               nvalid := spacedCompress(in, enc.spacedScratch, validBits, 
validBitsOffset)
+               enc.Put(enc.spacedScratch[:nvalid])
+               clear(enc.spacedScratch)
        } else {
                enc.Put(in)
        }
diff --git 
a/parquet/internal/encoding/delta_byte_array_spaced_benchmark_test.go 
b/parquet/internal/encoding/delta_byte_array_spaced_benchmark_test.go
new file mode 100644
index 00000000..ec85c627
--- /dev/null
+++ b/parquet/internal/encoding/delta_byte_array_spaced_benchmark_test.go
@@ -0,0 +1,84 @@
+// 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/bitutil"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/parquet"
+)
+
+func BenchmarkDeltaLengthByteArrayPutSpaced(b *testing.B) {
+       benchmarkDeltaByteArrayPutSpaced(b, 
parquet.Encodings.DeltaLengthByteArray)
+}
+
+func BenchmarkDeltaByteArrayPutSpaced(b *testing.B) {
+       benchmarkDeltaByteArrayPutSpaced(b, parquet.Encodings.DeltaByteArray)
+}
+
+func benchmarkDeltaByteArrayPutSpaced(b *testing.B, encoding parquet.Encoding) 
{
+       patterns := []struct {
+               name  string
+               valid func(int) bool
+       }{
+               {name: "all_valid", valid: func(int) bool { return true }},
+               {name: "ten_percent_null", valid: func(i int) bool { return 
i%10 != 0 }},
+               {name: "fifty_percent_null", valid: func(i int) bool { return 
i%2 != 0 }},
+               {name: "ninety_percent_null", valid: func(i int) bool { return 
i%10 == 0 }},
+       }
+
+       for _, length := range []int{1024, 64 * 1024} {
+               values := make([]parquet.ByteArray, length)
+               for i := range values {
+                       values[i] = 
parquet.ByteArray(fmt.Sprintf("partition/%06d", i))
+               }
+
+               for _, pattern := range patterns {
+                       b.Run(fmt.Sprintf("length_%d/%s", length, 
pattern.name), func(b *testing.B) {
+                               validBits := make([]byte, 
bitutil.BytesForBits(int64(length)))
+                               for i := range length {
+                                       if pattern.valid(i) {
+                                               bitutil.SetBit(validBits, i)
+                                       }
+                               }
+
+                               encoder := NewEncoder(parquet.Types.ByteArray, 
encoding, false, nil, memory.DefaultAllocator).(ByteArrayEncoder)
+                               defer encoder.Release()
+
+                               encode := func() {
+                                       encoder.PutSpaced(values, validBits, 0)
+                                       buf, err := encoder.FlushValues()
+                                       if err != nil {
+                                               b.Fatal(err)
+                                       }
+                                       buf.Release()
+                               }
+
+                               encode()
+                               b.SetBytes(int64(length * 16))
+                               b.ReportAllocs()
+                               b.ResetTimer()
+                               for b.Loop() {
+                                       encode()
+                               }
+                       })
+               }
+       }
+}
diff --git a/parquet/internal/encoding/delta_byte_array_spaced_test.go 
b/parquet/internal/encoding/delta_byte_array_spaced_test.go
new file mode 100644
index 00000000..6c4d0797
--- /dev/null
+++ b/parquet/internal/encoding/delta_byte_array_spaced_test.go
@@ -0,0 +1,135 @@
+// 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/bitutil"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/parquet"
+       "github.com/stretchr/testify/require"
+)
+
+func TestDeltaByteArrayPutSpacedReusesScratch(t *testing.T) {
+       const nvalues = deltaByteArrayBatchSize + 3
+
+       values := make([]parquet.ByteArray, nvalues)
+       validBits := make([]byte, bitutil.BytesForBits(nvalues))
+       for i := range values {
+               values[i] = parquet.ByteArray(fmt.Sprintf("value-%03d", i))
+               bitutil.SetBit(validBits, i)
+       }
+
+       tests := []struct {
+               name    string
+               new     func() ByteArrayEncoder
+               scratch func(ByteArrayEncoder) []parquet.ByteArray
+       }{
+               {
+                       name: "delta-length-byte-array",
+                       new: func() ByteArrayEncoder {
+                               return NewEncoder(parquet.Types.ByteArray, 
parquet.Encodings.DeltaLengthByteArray,
+                                       false, nil, 
memory.DefaultAllocator).(ByteArrayEncoder)
+                       },
+                       scratch: func(enc ByteArrayEncoder) []parquet.ByteArray 
{
+                               return 
enc.(*DeltaLengthByteArrayEncoder).spacedScratch
+                       },
+               },
+               {
+                       name: "delta-byte-array",
+                       new: func() ByteArrayEncoder {
+                               return NewEncoder(parquet.Types.ByteArray, 
parquet.Encodings.DeltaByteArray,
+                                       false, nil, 
memory.DefaultAllocator).(ByteArrayEncoder)
+                       },
+                       scratch: func(enc ByteArrayEncoder) []parquet.ByteArray 
{
+                               return 
enc.(*DeltaByteArrayEncoder).spacedScratch
+                       },
+               },
+       }
+
+       for _, tt := range tests {
+               t.Run(tt.name, func(t *testing.T) {
+                       enc := tt.new()
+                       defer enc.Release()
+
+                       enc.PutSpaced(values, validBits, 0)
+                       firstScratch := tt.scratch(enc)
+                       require.Len(t, firstScratch, nvalues)
+                       firstValue := &firstScratch[0]
+                       for i, value := range firstScratch {
+                               require.Nil(t, value, "scratch entry %d still 
references input data", i)
+                       }
+
+                       buf, err := enc.FlushValues()
+                       require.NoError(t, err)
+                       buf.Release()
+
+                       enc.PutSpaced(values[:1], validBits, 0)
+                       secondScratch := tt.scratch(enc)
+                       require.Len(t, secondScratch, 1)
+                       require.True(t, firstValue == &secondScratch[0], 
"scratch backing storage was not reused")
+                       for i, value := range 
secondScratch[:cap(secondScratch)] {
+                               require.Nil(t, value, "scratch entry %d still 
references input data", i)
+                       }
+
+                       buf, err = enc.FlushValues()
+                       require.NoError(t, err)
+                       buf.Release()
+               })
+       }
+}
+
+func TestDeltaByteArrayPutSpacedRoundTripWithOffset(t *testing.T) {
+       const nvalues = deltaByteArrayBatchSize*2 + 7
+       const validBitsOffset = int64(5)
+
+       values := make([]parquet.ByteArray, nvalues)
+       validBits := make([]byte, 
bitutil.BytesForBits(validBitsOffset+int64(nvalues)))
+       want := make([]parquet.ByteArray, 0, nvalues)
+       for i := range values {
+               values[i] = 
parquet.ByteArray(fmt.Sprintf("partition-%02d/value-%03d", i/11, i))
+               if i%7 != 2 {
+                       bitutil.SetBit(validBits, int(validBitsOffset)+i)
+                       want = append(want, values[i])
+               }
+       }
+
+       for _, encoding := range []parquet.Encoding{
+               parquet.Encodings.DeltaLengthByteArray,
+               parquet.Encodings.DeltaByteArray,
+       } {
+               t.Run(encoding.String(), func(t *testing.T) {
+                       enc := NewEncoder(parquet.Types.ByteArray, encoding, 
false, nil, memory.DefaultAllocator).(ByteArrayEncoder)
+                       defer enc.Release()
+
+                       enc.PutSpaced(values, validBits, validBitsOffset)
+                       buf, err := enc.FlushValues()
+                       require.NoError(t, err)
+                       defer buf.Release()
+
+                       dec := NewDecoder(parquet.Types.ByteArray, encoding, 
nil, memory.DefaultAllocator).(ByteArrayDecoder)
+                       require.NoError(t, dec.SetData(len(want), buf.Bytes()))
+                       got := make([]parquet.ByteArray, len(want))
+                       decoded, err := dec.Decode(got)
+                       require.NoError(t, err)
+                       require.Equal(t, len(want), decoded)
+                       require.Equal(t, want, got)
+               })
+       }
+}
diff --git a/parquet/internal/encoding/delta_length_byte_array.go 
b/parquet/internal/encoding/delta_length_byte_array.go
index 30ce53ff..4b74ff1d 100644
--- a/parquet/internal/encoding/delta_length_byte_array.go
+++ b/parquet/internal/encoding/delta_length_byte_array.go
@@ -39,6 +39,7 @@ type DeltaLengthByteArrayEncoder struct {
 
        lengthEncoder *DeltaBitPackInt32Encoder
        lengths       [deltaByteArrayBatchSize]int32
+       spacedScratch []parquet.ByteArray
 }
 
 // Put writes the provided slice of byte arrays to the encoder
@@ -62,9 +63,14 @@ func (enc *DeltaLengthByteArrayEncoder) Put(in 
[]parquet.ByteArray) {
 // accordingly before it is written to drop the null data from the write.
 func (enc *DeltaLengthByteArrayEncoder) PutSpaced(in []parquet.ByteArray, 
validBits []byte, validBitsOffset int64) {
        if validBits != nil {
-               data := make([]parquet.ByteArray, len(in))
-               nvalid := spacedCompress(in, data, validBits, validBitsOffset)
-               enc.Put(data[:nvalid])
+               if cap(enc.spacedScratch) < len(in) {
+                       enc.spacedScratch = make([]parquet.ByteArray, len(in))
+               } else {
+                       enc.spacedScratch = enc.spacedScratch[:len(in)]
+               }
+               nvalid := spacedCompress(in, enc.spacedScratch, validBits, 
validBitsOffset)
+               enc.Put(enc.spacedScratch[:nvalid])
+               clear(enc.spacedScratch)
        } else {
                enc.Put(in)
        }

Reply via email to