From 4609841cbceaa34e6abb7ba445db4c2d66946b00 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Sat, 29 Aug 2026 14:39:51 +0200 Subject: [PATCH] perf(parquet): reuse DELTA byte-array scratch --- parquet/internal/encoding/delta_byte_array.go | 12 +- .../delta_byte_array_spaced_benchmark_test.go | 84 +++++++++++ .../encoding/delta_byte_array_spaced_test.go | 135 ++++++++++++++++++ .../encoding/delta_length_byte_array.go | 12 +- 4 files changed, 237 insertions(+), 6 deletions(-) create mode 100644 parquet/internal/encoding/delta_byte_array_spaced_benchmark_test.go create mode 100644 parquet/internal/encoding/delta_byte_array_spaced_test.go diff --git a/parquet/internal/encoding/delta_byte_array.go b/parquet/internal/encoding/delta_byte_array.go index 86b7d5857..ae88d16a4 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 000000000..ec85c627e --- /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 000000000..6c4d0797b --- /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 30ce53ff5..4b74ff1dc 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) }