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 63833f78 perf(parquet): batch delta binary packed decoding (#1169)
63833f78 is described below
commit 63833f783b1e52cb0bb0dc0b4e72ad6a2077abff
Author: Minh Vu <[email protected]>
AuthorDate: Tue Aug 18 18:57:12 2026 +0200
perf(parquet): batch delta binary packed decoding (#1169)
### Rationale for this change
DELTA_BINARY_PACKED decoding called `BitReader.GetBatch` once for every
value in a miniblock. The miniblock is already the unit decoded by this
code, so this missed the existing batch unpackers.
### What changes are included in this PR?
- decode each miniblock into reusable scratch with one `GetBatch` call
- keep 33 to 64 bit values on the generic uint64 path instead of the 32
bit SIMD unpackers
- read SIMD input directly into pooled scratch so batched calls do not
add allocations
- add int32 and int64 benchmarks with nonzero packed widths
- add coverage for batched values from 33 through 64 bits
### Are these changes tested?
Yes. I ran the full Parquet test tree with the parquet-testing fixtures,
the internal utility and encoding tests with `noasm`, and a linux/amd64
cross-build.
Apple M1 Pro results for 65,536 values with `GOMAXPROCS=1` and 10
samples:
| Workload | Time |
| --- | ---: |
| int32 small deltas | -58.86% |
| int32 alternating wide deltas | -66.04% |
| int64 timestamp-like deltas | -59.72% |
| int64 deltas wider than 32 bits | -33.89% |
Allocations stay at 2 allocs/op in all four benchmarks.
---
parquet/internal/encoding/delta_bit_packing.go | 31 ++++---
.../encoding/delta_bit_packing_validation_test.go | 43 ++++++++++
.../internal/encoding/encoding_benchmarks_test.go | 96 +++++++++++++++++++++-
parquet/internal/encoding/encoding_test.go | 36 ++++++++
parquet/internal/utils/bit_packing_avx2_amd64.go | 11 +--
parquet/internal/utils/bit_packing_neon_arm64.go | 11 +--
parquet/internal/utils/bit_reader.go | 34 ++++----
parquet/internal/utils/bit_reader_test.go | 64 +++++++++++++++
8 files changed, 286 insertions(+), 40 deletions(-)
diff --git a/parquet/internal/encoding/delta_bit_packing.go
b/parquet/internal/encoding/delta_bit_packing.go
index 262a2d4e..e8b3a30f 100644
--- a/parquet/internal/encoding/delta_bit_packing.go
+++ b/parquet/internal/encoding/delta_bit_packing.go
@@ -20,6 +20,7 @@ import (
"bytes"
"errors"
"fmt"
+ "io"
"math"
"math/bits"
@@ -30,6 +31,8 @@ import (
"github.com/apache/arrow-go/v18/parquet/internal/utils"
)
+const deltaBitPackScratchSize = 1024
+
// see the deltaBitPack encoder for a description of the encoding format that
is
// used for delta-bitpacking.
type deltaBitPackDecoder[T int32 | int64] struct {
@@ -181,23 +184,31 @@ func (d *deltaBitPackDecoder[T]) unpackNextMini() error {
d.miniBlockValues = append(d.miniBlockValues,
T(d.lastVal))
}
} else {
- if cap(d.deltaBuf) < 1 {
- d.deltaBuf = make([]uint64, 1)
+ scratchSize := min(n, deltaBitPackScratchSize)
+ if cap(d.deltaBuf) < scratchSize {
+ d.deltaBuf = make([]uint64, scratchSize)
}
- d.deltaBuf = d.deltaBuf[:1]
+ d.deltaBuf = d.deltaBuf[:scratchSize]
minDelta := d.minDelta
- for j := 0; j < n; j++ {
- nread, err := d.bitdecoder.GetBatch(width, d.deltaBuf)
- if err != nil {
+ for remaining := n; remaining > 0; {
+ batchSize := min(remaining, len(d.deltaBuf))
+ nread, err := d.bitdecoder.GetBatch(width,
d.deltaBuf[:batchSize])
+ if nread != batchSize {
+ if err == nil || errors.Is(err, io.EOF) ||
errors.Is(err, io.ErrUnexpectedEOF) {
+ return io.ErrUnexpectedEOF
+ }
return err
}
- if nread != 1 {
- return errors.New("parquet: eof exception")
+ if err != nil {
+ return err
}
- d.lastVal += int64(d.deltaBuf[0]) + minDelta
- d.miniBlockValues = append(d.miniBlockValues,
T(d.lastVal))
+ for _, delta := range d.deltaBuf[:nread] {
+ d.lastVal += int64(delta) + minDelta
+ d.miniBlockValues = append(d.miniBlockValues,
T(d.lastVal))
+ }
+ remaining -= batchSize
}
}
d.miniBlockIdx++
diff --git a/parquet/internal/encoding/delta_bit_packing_validation_test.go
b/parquet/internal/encoding/delta_bit_packing_validation_test.go
index 90b4733c..030aa96a 100644
--- a/parquet/internal/encoding/delta_bit_packing_validation_test.go
+++ b/parquet/internal/encoding/delta_bit_packing_validation_test.go
@@ -17,6 +17,7 @@
package encoding
import (
+ "io"
"testing"
"github.com/apache/arrow-go/v18/arrow/memory"
@@ -74,3 +75,45 @@ func TestDeltaBitPackDecoderValidatesUsedBitWidth(t
*testing.T) {
})
}
}
+
+func TestDeltaBitPackDecoderRejectsTruncatedPackedMiniblock(t *testing.T) {
+ // Header: 1024 values/block, 1 miniblock, 1025 values, first value 0.
+ // The miniblock has a one-bit width, but only its first 32 packed
values are
+ // present. The missing values must not be reported as a clean EOF.
+ data := []byte{
+ 128, 8, // block size
+ 1, // miniblocks per block
+ 129, 8, // total values
+ 0, // first value
+ 0, // minimum delta
+ 1, // first miniblock bit width
+ 0, 0, 0, 0,
+ }
+
+ dec := NewDecoder(parquet.Types.Int64,
parquet.Encodings.DeltaBinaryPacked, nil, memory.DefaultAllocator)
+ require.NoError(t, dec.SetData(1025, data))
+
+ _, err := dec.Discard(1025)
+ require.ErrorIs(t, err, io.ErrUnexpectedEOF)
+}
+
+func TestDeltaBitPackDecoderBoundsPackedScratch(t *testing.T) {
+ // A malformed header can declare a large miniblock without providing
any
+ // packed values. The reusable packed-value scratch must stay bounded
instead
+ // of growing to the attacker-controlled miniblock size.
+ data := []byte{
+ 128, 32, // block size: 4096 values
+ 1, // miniblocks per block
+ 129, 32, // total values: first value plus one block
+ 0, // first value
+ 0, // minimum delta
+ 1, // first miniblock bit width
+ }
+
+ dec := NewDecoder(parquet.Types.Int64,
parquet.Encodings.DeltaBinaryPacked, nil,
memory.DefaultAllocator).(*deltaBitPackDecoder[int64])
+ require.NoError(t, dec.SetData(4097, data))
+
+ _, err := dec.Discard(2)
+ require.Error(t, err)
+ require.LessOrEqual(t, cap(dec.deltaBuf), deltaBitPackScratchSize)
+}
diff --git a/parquet/internal/encoding/encoding_benchmarks_test.go
b/parquet/internal/encoding/encoding_benchmarks_test.go
index 09c176d9..781967a9 100644
--- a/parquet/internal/encoding/encoding_benchmarks_test.go
+++ b/parquet/internal/encoding/encoding_benchmarks_test.go
@@ -735,13 +735,49 @@ func BenchmarkDeltaBinaryPackedEncodingInt32(b
*testing.B) {
}
func BenchmarkDeltaBinaryPackedDecodingInt32(b *testing.B) {
- for sz := MINSIZE; sz < MAXSIZE+1; sz *= 2 {
- b.Run(fmt.Sprintf("len %d", sz), func(b *testing.B) {
- output := make([]int32, sz)
- values := make([]int32, sz)
+ patterns := []struct {
+ name string
+ values func(int) []int32
+ }{
+ {"constant", func(size int) []int32 {
+ values := make([]int32, size)
for idx := range values {
values[idx] = 64
}
+ return values
+ }},
+ {"varying-small-deltas", func(size int) []int32 {
+ values := make([]int32, size)
+ for idx := 1; idx < size; idx++ {
+ values[idx] = values[idx-1] + int32(idx%8)
+ }
+ return values
+ }},
+ {"alternating-wide-deltas", func(size int) []int32 {
+ values := make([]int32, size)
+ for idx := 1; idx < size; idx++ {
+ if idx%2 == 0 {
+ values[idx] = values[idx-1] + 1_000_000
+ } else {
+ values[idx] = values[idx-1] - 1_000_000
+ }
+ }
+ return values
+ }},
+ }
+
+ for _, pattern := range patterns {
+ b.Run(pattern.name, func(b *testing.B) {
+ benchmarkDeltaBinaryPackedDecodingInt32(b,
pattern.values)
+ })
+ }
+}
+
+func benchmarkDeltaBinaryPackedDecodingInt32(b *testing.B, makeValues
func(int) []int32) {
+ for sz := MINSIZE; sz < MAXSIZE+1; sz *= 2 {
+ b.Run(fmt.Sprintf("len-%d", sz), func(b *testing.B) {
+ output := make([]int32, sz)
+ values := makeValues(sz)
encoder := encoding.NewEncoder(parquet.Types.Int32,
parquet.Encodings.DeltaBinaryPacked,
false, nil,
memory.DefaultAllocator).(encoding.Int32Encoder)
encoder.Put(values)
@@ -751,6 +787,7 @@ func BenchmarkDeltaBinaryPackedDecodingInt32(b *testing.B) {
decoder := encoding.NewDecoder(parquet.Types.Int32,
parquet.Encodings.DeltaBinaryPacked, nil, memory.DefaultAllocator)
b.ResetTimer()
b.SetBytes(int64(len(values) * arrow.Int32SizeBytes))
+ b.ReportAllocs()
for n := 0; n < b.N; n++ {
decoder.SetData(sz, buf.Bytes())
decoder.(encoding.Int32Decoder).Decode(output)
@@ -759,6 +796,57 @@ func BenchmarkDeltaBinaryPackedDecodingInt32(b *testing.B)
{
}
}
+func BenchmarkDeltaBinaryPackedDecodingInt64(b *testing.B) {
+ patterns := []struct {
+ name string
+ values func(int) []int64
+ }{
+ {"varying-timestamp-deltas", func(size int) []int64 {
+ values := make([]int64, size)
+ for idx := 1; idx < size; idx++ {
+ values[idx] = values[idx-1] + 1_000_000 +
int64(idx%1024)
+ }
+ return values
+ }},
+ {"alternating-wide-deltas", func(size int) []int64 {
+ values := make([]int64, size)
+ for idx := 1; idx < size; idx++ {
+ if idx%2 == 0 {
+ values[idx] = values[idx-1] +
1_000_000_000_000
+ } else {
+ values[idx] = values[idx-1] -
1_000_000_000_000
+ }
+ }
+ return values
+ }},
+ }
+
+ for _, pattern := range patterns {
+ b.Run(pattern.name, func(b *testing.B) {
+ for sz := MINSIZE; sz < MAXSIZE+1; sz *= 2 {
+ b.Run(fmt.Sprintf("len-%d", sz), func(b
*testing.B) {
+ output := make([]int64, sz)
+ values := pattern.values(sz)
+ encoder :=
encoding.NewEncoder(parquet.Types.Int64, parquet.Encodings.DeltaBinaryPacked,
+ false, nil,
memory.DefaultAllocator).(encoding.Int64Encoder)
+ encoder.Put(values)
+ buf, _ := encoder.FlushValues()
+ defer buf.Release()
+
+ decoder :=
encoding.NewDecoder(parquet.Types.Int64, parquet.Encodings.DeltaBinaryPacked,
nil, memory.DefaultAllocator)
+ b.ResetTimer()
+ b.SetBytes(int64(len(values) *
arrow.Int64SizeBytes))
+ b.ReportAllocs()
+ for n := 0; n < b.N; n++ {
+ decoder.SetData(sz, buf.Bytes())
+
decoder.(encoding.Int64Decoder).Decode(output)
+ }
+ })
+ }
+ })
+ }
+}
+
// Extended MemoTable benchmarks for int64
func BenchmarkMemoTableInt64(b *testing.B) {
tests := []struct {
diff --git a/parquet/internal/encoding/encoding_test.go
b/parquet/internal/encoding/encoding_test.go
index 7a73df08..3afac422 100644
--- a/parquet/internal/encoding/encoding_test.go
+++ b/parquet/internal/encoding/encoding_test.go
@@ -956,6 +956,42 @@ func TestWriteDeltaBitPackedInt64(t *testing.T) {
})
}
+func TestDeltaBitPackedInt64DecodeAfterDiscard(t *testing.T) {
+ column := schema.NewColumn(schema.NewInt64Node("int64",
parquet.Repetitions.Required, -1), 0, 0)
+ values := make([]int64, 257)
+ for idx := 1; idx < len(values); idx++ {
+ delta := int64(1<<40) + int64(idx%17)
+ if idx%2 != 0 {
+ delta = -delta
+ }
+ values[idx] = values[idx-1] + delta
+ }
+
+ enc := encoding.NewEncoder(parquet.Types.Int64,
parquet.Encodings.DeltaBinaryPacked, false, column, memory.DefaultAllocator)
+ enc.(encoding.Int64Encoder).Put(values)
+ buf, err := enc.FlushValues()
+ require.NoError(t, err)
+ defer buf.Release()
+
+ for _, discard := range []int{0, 1, 31, 32, 33, 127, 128, 129, 256} {
+ t.Run(fmt.Sprintf("discard=%d", discard), func(t *testing.T) {
+ dec := encoding.NewDecoder(parquet.Types.Int64,
parquet.Encodings.DeltaBinaryPacked, column, memory.DefaultAllocator)
+ int64Dec := dec.(encoding.Int64Decoder)
+ require.NoError(t, int64Dec.SetData(len(values),
buf.Bytes()))
+
+ n, err := int64Dec.Discard(discard)
+ require.NoError(t, err)
+ require.Equal(t, discard, n)
+
+ out := make([]int64, len(values)-discard)
+ n, err = int64Dec.Decode(out)
+ require.NoError(t, err)
+ require.Equal(t, len(out), n)
+ assert.Equal(t, values[discard:], out)
+ })
+ }
+}
+
func TestDeltaLengthByteArrayEncoding(t *testing.T) {
column := schema.NewColumn(schema.NewByteArrayNode("bytearray",
parquet.Repetitions.Required, -1), 0, 0)
diff --git a/parquet/internal/utils/bit_packing_avx2_amd64.go
b/parquet/internal/utils/bit_packing_avx2_amd64.go
index 64fa86d9..74eb936b 100644
--- a/parquet/internal/utils/bit_packing_avx2_amd64.go
+++ b/parquet/internal/utils/bit_packing_avx2_amd64.go
@@ -42,17 +42,18 @@ func unpack32Avx2(in io.Reader, out []uint32, nbits int)
(int, error) {
defer bufferPool.Put(buffer)
buffer.Reset()
buffer.Grow(n)
- nread, err := io.CopyN(buffer, in, int64(n))
- if err == io.EOF && int(nread)%(nbits*4) != 0 {
- err = io.ErrUnexpectedEOF
+ packed := buffer.AvailableBuffer()[:n]
+ nread, err := io.ReadFull(in, packed)
+ if err == io.ErrUnexpectedEOF && nread%(nbits*4) == 0 {
+ err = io.EOF
}
- completeBatch := int(nread) * 8 / nbits / 32 * 32
+ completeBatch := nread * 8 / nbits / 32 * 32
if completeBatch == 0 {
return 0, err
}
var (
- input = unsafe.Pointer(&buffer.Bytes()[0])
+ input = unsafe.Pointer(&packed[0])
output = unsafe.Pointer(&out[0])
)
diff --git a/parquet/internal/utils/bit_packing_neon_arm64.go
b/parquet/internal/utils/bit_packing_neon_arm64.go
index e474f64c..fa204657 100644
--- a/parquet/internal/utils/bit_packing_neon_arm64.go
+++ b/parquet/internal/utils/bit_packing_neon_arm64.go
@@ -42,17 +42,18 @@ func unpack32NEON(in io.Reader, out []uint32, nbits int)
(int, error) {
defer bufferPool.Put(buffer)
buffer.Reset()
buffer.Grow(n)
- nread, err := io.CopyN(buffer, in, int64(n))
- if err == io.EOF && int(nread)%(nbits*4) != 0 {
- err = io.ErrUnexpectedEOF
+ packed := buffer.AvailableBuffer()[:n]
+ nread, err := io.ReadFull(in, packed)
+ if err == io.ErrUnexpectedEOF && nread%(nbits*4) == 0 {
+ err = io.EOF
}
- completeBatch := int(nread) * 8 / nbits / 32 * 32
+ completeBatch := nread * 8 / nbits / 32 * 32
if completeBatch == 0 {
return 0, err
}
var (
- input = unsafe.Pointer(&buffer.Bytes()[0])
+ input = unsafe.Pointer(&packed[0])
output = unsafe.Pointer(&out[0])
)
diff --git a/parquet/internal/utils/bit_reader.go
b/parquet/internal/utils/bit_reader.go
index 862394cc..39575d0c 100644
--- a/parquet/internal/utils/bit_reader.go
+++ b/parquet/internal/utils/bit_reader.go
@@ -555,24 +555,26 @@ func (b *BitReader) GetBatch(bits uint, out []uint64)
(int, error) {
}
}
- if _, err := b.reader.Seek(b.byteoffset, io.SeekStart); err != nil {
- return i, err
- }
- for i < length {
- // unpack groups of 32 bytes at a time into a buffer since it's
more efficient
- unpackSize := utils.Min(buflen, length-i)
- numUnpacked, err := unpack32(b.reader,
b.unpackBuf[:unpackSize], int(bits))
-
- for k := 0; k < numUnpacked; k++ {
- out[i+k] = uint64(b.unpackBuf[k])
- }
- i += numUnpacked
- b.byteoffset += int64(numUnpacked * int(bits) / 8)
- if err != nil {
+ if bits <= 32 {
+ if _, err := b.reader.Seek(b.byteoffset, io.SeekStart); err !=
nil {
return i, err
}
- if numUnpacked == 0 {
- break
+ for i < length {
+ // unpack groups of 32 bytes at a time into a buffer
since it's more efficient
+ unpackSize := utils.Min(buflen, length-i)
+ numUnpacked, err := unpack32(b.reader,
b.unpackBuf[:unpackSize], int(bits))
+
+ for k := 0; k < numUnpacked; k++ {
+ out[i+k] = uint64(b.unpackBuf[k])
+ }
+ i += numUnpacked
+ b.byteoffset += int64(numUnpacked * int(bits) / 8)
+ if err != nil {
+ return i, err
+ }
+ if numUnpacked == 0 {
+ break
+ }
}
}
diff --git a/parquet/internal/utils/bit_reader_test.go
b/parquet/internal/utils/bit_reader_test.go
index 502bbbfb..0468ff13 100644
--- a/parquet/internal/utils/bit_reader_test.go
+++ b/parquet/internal/utils/bit_reader_test.go
@@ -360,6 +360,70 @@ func TestBitArrayVals(t *testing.T) {
}
}
+func TestBitReaderGetBatchWideValues(t *testing.T) {
+ for width := uint(33); width <= 64; width++ {
+ t.Run(fmt.Sprintf("width=%d", width), func(t *testing.T) {
+ const nvalues = 32
+ values := make([]uint64, nvalues)
+ mask := uint64(math.MaxUint64)
+ if width < 64 {
+ mask = 1<<width - 1
+ }
+ for idx := range values {
+ values[idx] = uint64(idx) * 0x9e3779b97f4a7c15
& mask
+ }
+
+ buf := make([]byte,
bitutil.BytesForBits(int64(width*nvalues)))
+ writer :=
utils.NewBitWriter(utils.NewWriterAtBuffer(buf))
+ for _, value := range values {
+ assert.NoError(t, writer.WriteValue(value,
width))
+ }
+ writer.Flush(false)
+
+ reader := utils.NewBitReader(bytes.NewReader(buf))
+ actual := make([]uint64, nvalues)
+ n, err := reader.GetBatch(width, actual)
+ assert.NoError(t, err)
+ assert.Equal(t, nvalues, n)
+ assert.Equal(t, values, actual)
+ })
+ }
+}
+
+func TestBitReaderGetBatchKeepsFollowingValues(t *testing.T) {
+ for width := uint(1); width <= 64; width++ {
+ t.Run(fmt.Sprintf("width=%d", width), func(t *testing.T) {
+ const nvalues = 65
+ values := make([]uint64, nvalues)
+ mask := uint64(math.MaxUint64)
+ if width < 64 {
+ mask = 1<<width - 1
+ }
+ for idx := range values {
+ values[idx] = (uint64(idx)*0x9e3779b97f4a7c15 +
uint64(idx/3)) & mask
+ }
+
+ buf := make([]byte,
bitutil.BytesForBits(int64(width*nvalues)))
+ writer :=
utils.NewBitWriter(utils.NewWriterAtBuffer(buf))
+ for _, value := range values {
+ assert.NoError(t, writer.WriteValue(value,
width))
+ }
+ writer.Flush(false)
+
+ reader := utils.NewBitReader(bytes.NewReader(buf))
+ actual := make([]uint64, nvalues-1)
+ n, err := reader.GetBatch(width, actual)
+ assert.NoError(t, err)
+ assert.Equal(t, nvalues-1, n)
+ assert.Equal(t, values[:nvalues-1], actual)
+
+ value, ok := reader.GetValue(int(width))
+ assert.True(t, ok)
+ assert.Equal(t, values[nvalues-1], value)
+ })
+ }
+}
+
func TestBitReaderRejectsTruncatedPackedBatches(t *testing.T) {
for width := 1; width <= 32; width++ {
for _, batchSize := range []int{32, 64} {