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)
}