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 7573b714 perf(arrow/array): avoid temporary buffers for binary
concatenation (#1215)
7573b714 is described below
commit 7573b714647cf0a43e40bbafc62ff7ffec33cf74
Author: Minh Vu <[email protected]>
AuthorDate: Fri Aug 28 22:09:09 2026 +0200
perf(arrow/array): avoid temporary buffers for binary concatenation (#1215)
## Summary
- copy binary and string values directly from source buffers
- avoid creating temporary `memory.Buffer` wrappers for each input chunk
- keep offset rebasing and overflow checks in the direct path
- add sliced binary/string coverage and a variable-width benchmark
## Benchmarks
Apple M1 Pro, 65,536 values, 8-byte values:
| chunks | old allocs | new allocs | old B/op | new B/op |
| ---: | ---: | ---: | ---: | ---: |
| 1 | 12 | 8 | 803,411 | 803,229 |
| 1,024 | 2,059 | 9 | 1,020,824 | 838,035 |
| 8,192 | 16,395 | 9 | 2,507,159 | 1,065,358 |
## Tests
- `PARQUET_TEST_DATA="$PWD/parquet-testing/data" go test -count=1 ./...`
- `go vet ./arrow/array`
---
arrow/array/concat.go | 120 ++++++++++++++++++++++++++++++++++++--
arrow/array/concat_test.go | 140 +++++++++++++++++++++++++++++++++++++++++++++
2 files changed, 256 insertions(+), 4 deletions(-)
diff --git a/arrow/array/concat.go b/arrow/array/concat.go
index 510ebe54..9a35daae 100644
--- a/arrow/array/concat.go
+++ b/arrow/array/concat.go
@@ -190,6 +190,121 @@ func concatFixedWidthBuffers(data []arrow.ArrayData, idx,
byteWidth, length int,
return out
}
+func concatBinaryBuffers(data []arrow.ArrayData, byteWidth, length int, out
*Data, mem memory.Allocator) error {
+ offsetBuffer := memory.NewResizableBuffer(mem)
+ out.buffers[1] = offsetBuffer
+ offsetBuffer.Resize(byteWidth * (length + 1))
+
+ var (
+ valueRanges []rng
+ err error
+ )
+ switch byteWidth {
+ case arrow.Int64SizeBytes:
+ valueRanges, err = handle64BitOffsetsData(data, offsetBuffer,
length)
+ default:
+ valueRanges, err = handle32BitOffsetsData(data, offsetBuffer,
length)
+ }
+ if err != nil {
+ return err
+ }
+
+ valueLength := 0
+ for _, r := range valueRanges {
+ valueLength += r.len
+ }
+
+ valueBuffer := memory.NewResizableBuffer(mem)
+ out.buffers[2] = valueBuffer
+ valueBuffer.Resize(valueLength)
+ dst := valueBuffer.Bytes()
+ for i, d := range data {
+ r := valueRanges[i]
+ if r.len == 0 {
+ continue
+ }
+
+ buf := d.Buffers()[2]
+ copy(dst, buf.Bytes()[r.offset:r.offset+r.len])
+ dst = dst[r.len:]
+ }
+ return nil
+}
+
+func handle32BitOffsetsData(data []arrow.ArrayData, out *memory.Buffer, outLen
int) ([]rng, error) {
+ dst := arrow.Int32Traits.CastFromBytes(out.Bytes())
+ valueRanges := make([]rng, len(data))
+ nextOffset := int32(0)
+ nextElem := 0
+ for i, d := range data {
+ if d.Len() == 0 {
+ continue
+ }
+
+ buf := d.Buffers()[1]
+ if buf == nil {
+ return nil, errors.New("array/concat: binary array is
missing an offset buffer")
+ }
+ src := arrow.Int32Traits.CastFromBytes(buf.Bytes())
+ begin := d.Offset()
+ end := begin + d.Len()
+ startOffset, endOffset := src[begin], src[end]
+ valueLength := int(endOffset) - int(startOffset)
+
+ if valueLength < 0 || int64(nextOffset)+int64(valueLength) >
math.MaxInt32 {
+ return nil, errors.New("offset overflow while
concatenating arrays")
+ }
+
+ valueRanges[i] = rng{offset: int(startOffset), len: valueLength}
+ adj := nextOffset - startOffset
+ for j, o := range src[begin:end] {
+ dst[nextElem+j] = adj + o
+ }
+ nextElem += d.Len()
+ nextOffset += int32(valueLength)
+ }
+
+ dst[outLen] = nextOffset
+ return valueRanges, nil
+}
+
+func handle64BitOffsetsData(data []arrow.ArrayData, out *memory.Buffer, outLen
int) ([]rng, error) {
+ dst := arrow.Int64Traits.CastFromBytes(out.Bytes())
+ valueRanges := make([]rng, len(data))
+ nextOffset := int64(0)
+ nextElem := 0
+ for i, d := range data {
+ if d.Len() == 0 {
+ continue
+ }
+
+ buf := d.Buffers()[1]
+ if buf == nil {
+ return nil, errors.New("array/concat: binary array is
missing an offset buffer")
+ }
+ src := arrow.Int64Traits.CastFromBytes(buf.Bytes())
+ begin := d.Offset()
+ end := begin + d.Len()
+ startOffset, endOffset := src[begin], src[end]
+ valueLength := int(endOffset) - int(startOffset)
+
+ if valueLength < 0 || nextOffset >
math.MaxInt64-int64(valueLength) {
+ return nil, errors.New("offset overflow while
concatenating arrays")
+ }
+
+ valueRanges[i] = rng{offset: int(startOffset), len: valueLength}
+ adj := nextOffset - startOffset
+ for j, o := range src[begin:end] {
+ dst[nextElem+j] = adj + o
+ }
+ nextElem += d.Len()
+ nextOffset += int64(valueLength)
+ }
+
+ dst[outLen] = nextOffset
+ return valueRanges, nil
+}
+
func handle32BitOffsets(outLen int, buffers []*memory.Buffer, out
*memory.Buffer) (*memory.Buffer, []rng, error) {
dst := arrow.Int32Traits.CastFromBytes(out.Bytes())
valuesRanges := make([]rng, len(buffers))
@@ -647,12 +762,9 @@ func concat(data []arrow.ArrayData, mem memory.Allocator)
(arr arrow.ArrayData,
}
case arrow.BinaryDataType:
offsetWidth := dt.Layout().Buffers[1].ByteWidth
- offsetBuffer, valueRanges, err :=
concatOffsets(gatherFixedBuffers(data, 1, offsetWidth), offsetWidth, mem)
- if err != nil {
+ if err := concatBinaryBuffers(data, offsetWidth, out.length,
out, mem); err != nil {
return nil, err
}
- out.buffers[1] = offsetBuffer
- out.buffers[2] = concatBuffers(gatherBufferRanges(data, 2,
valueRanges), mem)
case *arrow.ListType:
offsetWidth := dt.Layout().Buffers[1].ByteWidth
offsetBuffer, valueRanges, err :=
concatOffsets(gatherFixedBuffers(data, 1, offsetWidth), offsetWidth, mem)
diff --git a/arrow/array/concat_test.go b/arrow/array/concat_test.go
index df906bd9..135f9991 100644
--- a/arrow/array/concat_test.go
+++ b/arrow/array/concat_test.go
@@ -190,6 +190,146 @@ func TestConcatenateFixedWidthSlices(t *testing.T) {
}
}
+func BenchmarkConcatenateBinary(b *testing.B) {
+ const totalValues = 1 << 16
+
+ types := []struct {
+ name string
+ dt arrow.BinaryDataType
+ }{
+ {"binary", arrow.BinaryTypes.Binary},
+ {"string", arrow.BinaryTypes.String},
+ {"large_binary", arrow.BinaryTypes.LargeBinary},
+ {"large_string", arrow.BinaryTypes.LargeString},
+ }
+
+ for _, tt := range types {
+ for _, valueSize := range []int{8, 32} {
+ tt, valueSize := tt, valueSize
+ b.Run(fmt.Sprintf("%s/value_size=%d", tt.name,
valueSize), func(b *testing.B) {
+ mem := memory.NewGoAllocator()
+ backing := makeConcatenateBinaryArray(mem,
tt.dt, totalValues, valueSize)
+ defer backing.Release()
+
+ for _, chunkCount := range []int{1, 8, 64,
1024, 8192} {
+ chunkCount := chunkCount
+ b.Run(fmt.Sprintf("chunks=%d",
chunkCount), func(b *testing.B) {
+ chunkSize := totalValues /
chunkCount
+ inputs := make([]arrow.Array,
chunkCount)
+ for i := range inputs {
+ begin := int64(i *
chunkSize)
+ inputs[i] =
array.NewSlice(backing, begin, begin+int64(chunkSize))
+ }
+ defer func() {
+ for _, input := range
inputs {
+ input.Release()
+ }
+ }()
+
+ b.SetBytes(int64(totalValues *
valueSize))
+ b.ReportAllocs()
+ b.ResetTimer()
+ for i := 0; i < b.N; i++ {
+ result, err :=
array.Concatenate(inputs, mem)
+ if err != nil {
+ b.Fatal(err)
+ }
+ if result.Len() !=
totalValues {
+
b.Fatalf("result length = %d, want %d", result.Len(), totalValues)
+ }
+ result.Release()
+ }
+ })
+ }
+ })
+ }
+ }
+}
+
+func TestConcatenateBinarySlices(t *testing.T) {
+ types := []struct {
+ name string
+ dt arrow.BinaryDataType
+ }{
+ {"binary", arrow.BinaryTypes.Binary},
+ {"string", arrow.BinaryTypes.String},
+ {"large_binary", arrow.BinaryTypes.LargeBinary},
+ {"large_string", arrow.BinaryTypes.LargeString},
+ }
+
+ for _, tt := range types {
+ t.Run(tt.name, func(t *testing.T) {
+ mem :=
memory.NewCheckedAllocator(memory.DefaultAllocator)
+ defer mem.AssertSize(t, 0)
+
+ backing := makeConcatenateBinaryArrayWithValidity(mem,
tt.dt,
+ []string{"zero", "one", "two", "", "four"},
+ []bool{true, false, true, true, true})
+ defer backing.Release()
+
+ inputs := []arrow.Array{
+ array.NewSlice(backing, 1, 4),
+ array.NewSlice(backing, 2, 2),
+ array.NewSlice(backing, 4, 5),
+ array.NewSlice(backing, 0, 1),
+ }
+ for _, input := range inputs {
+ defer input.Release()
+ }
+
+ actual, err := array.Concatenate(inputs, mem)
+ require.NoError(t, err)
+ defer actual.Release()
+
+ expected := makeConcatenateBinaryArrayWithValidity(mem,
tt.dt,
+ []string{"one", "two", "", "four", "zero"},
+ []bool{false, true, true, true, true})
+ defer expected.Release()
+
+ assert.True(t, array.Equal(expected, actual))
+ assert.Equal(t, 1, actual.NullN())
+ switch actual.DataType().Layout().Buffers[1].ByteWidth {
+ case arrow.Int32SizeBytes:
+ assert.Equal(t, []int32{0, 0, 3, 3, 7, 11},
arrow.Int32Traits.CastFromBytes(actual.Data().Buffers()[1].Bytes()))
+ case arrow.Int64SizeBytes:
+ assert.Equal(t, []int64{0, 0, 3, 3, 7, 11},
arrow.Int64Traits.CastFromBytes(actual.Data().Buffers()[1].Bytes()))
+ }
+ })
+ }
+}
+
+func makeConcatenateBinaryArray(mem memory.Allocator, dt arrow.BinaryDataType,
length, valueSize int) arrow.Array {
+ builder := array.NewBinaryBuilder(mem, dt)
+ builder.Reserve(length)
+ builder.ReserveData(length * valueSize)
+ value := strings.Repeat("a", valueSize)
+ for i := 0; i < length; i++ {
+ builder.AppendString(value)
+ }
+ return finishConcatenateBinaryArray(builder)
+}
+
+func makeConcatenateBinaryArrayWithValidity(mem memory.Allocator, dt
arrow.BinaryDataType, values []string, valid []bool) arrow.Array {
+ builder := array.NewBinaryBuilder(mem, dt)
+ builder.Reserve(len(values))
+ for i, value := range values {
+ if !valid[i] {
+ builder.AppendNull()
+ continue
+ }
+ builder.AppendString(value)
+ }
+ return finishConcatenateBinaryArray(builder)
+}
+
+func finishConcatenateBinaryArray(builder *array.BinaryBuilder) arrow.Array {
+ raw := builder.NewArray()
+ builder.Release()
+ result := array.MakeFromData(raw.Data())
+ raw.Release()
+ return result
+}
+
type ConcatTestSuite struct {
suite.Suite