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 b3e5d9c3 perf(parquet/pqarrow): use synchronous path for serial reads
(#1235)
b3e5d9c3 is described below
commit b3e5d9c37e13bc0228d52bdf25a132e6df843b2b
Author: Minh Vu <[email protected]>
AuthorDate: Mon Aug 31 21:29:04 2026 +0200
perf(parquet/pqarrow): use synchronous path for serial reads (#1235)
## Summary
- `ArrowReadProperties.Parallel` is false by default.
- `GetFieldReaders` now uses a direct loop for serial reads.
- `ReadRowGroups` now reads columns directly for serial reads.
- The parallel path is unchanged.
- Added a benchmark for 1, 8, 32, and 128 columns with small row groups.
## Benchmark
Apple M1 Pro. Median of 5 runs. The benchmark uses `Parallel=false`.
| columns | upstream main | this PR | time change | allocs/op |
| ---: | ---: | ---: | ---: | ---: |
| 1 | 156.9 us | 128.8 us | -18% | 1124 -> 1106 |
| 8 | 1.05 ms | 967 us | -8% | 7848 -> 7810 |
| 32 | 3.38 ms | 3.24 ms | -4% | 31413 -> 31296 |
| 128 | 11.92 ms | 11.18 ms | -6% | 124896 -> 124436 |
Command:
```text
go test ./parquet/pqarrow -run '^$' -bench '^BenchmarkReadTableSerial$'
-benchmem -benchtime=1s -count=5
```
## Tests
- `go test ./parquet/pqarrow -count=1`
- `go test -race ./parquet/pqarrow -run
'^(TestGetFieldReadersReleasesPartialReadersOnError|TestRecordReaderSerial|TestRecordReaderSeekToRow|TestRecordReaderParallel)$'
-count=1`
- `go vet ./parquet/pqarrow`
---
parquet/pqarrow/file_reader.go | 75 +++++++++-
parquet/pqarrow/file_reader_bench_test.go | 104 ++++++++++++++
parquet/pqarrow/file_reader_extension_test.go | 190 ++++++++++++++++++++++++++
3 files changed, 364 insertions(+), 5 deletions(-)
diff --git a/parquet/pqarrow/file_reader.go b/parquet/pqarrow/file_reader.go
index ae659d22..36382a6a 100644
--- a/parquet/pqarrow/file_reader.go
+++ b/parquet/pqarrow/file_reader.go
@@ -253,6 +253,24 @@ func (fr *FileReader) GetFieldReaders(ctx context.Context,
colIndices, rowGroups
out := make([]*ColumnReader, len(fieldIndices))
outFields := make([]arrow.Field, len(fieldIndices))
+ if !fr.Props.Parallel {
+ for idx, fidx := range fieldIndices {
+ rdr, err := fr.GetFieldReader(ctx, fidx,
includedLeaves, rowGroups)
+ if err != nil {
+ for _, rdr := range out {
+ if rdr != nil {
+ rdr.Release()
+ }
+ }
+ return nil, nil, err
+ }
+ outFields[idx] = *rdr.Field()
+ out[idx] = rdr
+ }
+
+ return out, arrow.NewSchema(outFields, fr.Manifest.SchemaMeta),
nil
+ }
+
// Load batches in parallel
// When reading structs with large numbers of columns, the serial load
is very slow.
// This is especially true when reading Cloud Storage. Loading
concurrently
@@ -260,9 +278,6 @@ func (fr *FileReader) GetFieldReaders(ctx context.Context,
colIndices, rowGroups
// GetFieldReader causes read operations, when issued serially on large
numbers of columns,
// this is super time consuming. Get field readers concurrently.
g, gctx := errgroup.WithContext(ctx)
- if !fr.Props.Parallel {
- g.SetLimit(1)
- }
for idx, fidx := range fieldIndices {
idx, fidx := idx, fidx // create concurrent copy
g.Go(func() error {
@@ -305,6 +320,15 @@ func (fr *FileReader) ReadColumn(rowGroups []int, rdr
*ColumnReader) (*arrow.Chu
return rdr.NextBatch(recs)
}
+func (fr *FileReader) readColumn(rowGroups []int, rdr *ColumnReader) (data
*arrow.Chunked, err error) {
+ defer func() {
+ if pErr := recover(); pErr != nil {
+ err = utils.FormatRecoveredError("panic while reading",
pErr)
+ }
+ }()
+ return fr.ReadColumn(rowGroups, rdr)
+}
+
// ReadTable reads the entire file into an array.Table
func (fr *FileReader) ReadTable(ctx context.Context) (arrow.Table, error) {
var (
@@ -370,6 +394,42 @@ func (fr *FileReader) ReadRowGroups(ctx context.Context,
indices, rowGroups []in
return nil, err
}
+ if !fr.Props.Parallel {
+ columns := make([]arrow.Column, sc.NumFields())
+ defer releaseColumns(columns)
+ defer func() {
+ for _, rdr := range readers {
+ rdr.Release()
+ }
+ }()
+
+ for idx, rdr := range readers {
+ if err := ctx.Err(); err != nil {
+ return nil, err
+ }
+
+ data, err := fr.readColumn(rowGroups, rdr)
+ if err != nil {
+ if data != nil {
+ data.Release()
+ }
+ return nil, err
+ }
+ columns[idx] = *arrow.NewColumn(sc.Field(idx), data)
+ data.Release()
+ }
+ if err := ctx.Err(); err != nil {
+ return nil, err
+ }
+
+ var nrows int
+ if len(columns) > 0 {
+ nrows = columns[0].Len()
+ }
+
+ return array.NewTable(sc, columns, int64(nrows)), nil
+ }
+
// producer-consumer parallelization
var (
np = 1
@@ -403,7 +463,7 @@ func (fr *FileReader) ReadRowGroups(ctx context.Context,
indices, rowGroups []in
return
}
- chnked, err := fr.ReadColumn(rowGroups,
r.rdr)
+ chnked, err := fr.readColumn(rowGroups,
r.rdr)
// pass the result column data to the
result channel
// for the consumer goroutine to process
results <- resultPair{r.idx, chnked,
err}
@@ -432,6 +492,9 @@ func (fr *FileReader) ReadRowGroups(ctx context.Context,
indices, rowGroups []in
defer releaseColumns(columns)
for data := range results {
if data.err != nil {
+ if data.data != nil {
+ data.data.Release()
+ }
err = data.err
cancel()
break
@@ -448,7 +511,9 @@ func (fr *FileReader) ReadRowGroups(ctx context.Context,
indices, rowGroups []in
// so the goroutines don't leak and so memory can get cleaned
up. we already
// cancelled the context, so we're just consuming anything that
was already queued up.
for data := range results {
- data.data.Release()
+ if data.data != nil {
+ data.data.Release()
+ }
}
return nil, err
}
diff --git a/parquet/pqarrow/file_reader_bench_test.go
b/parquet/pqarrow/file_reader_bench_test.go
new file mode 100644
index 00000000..84578c0c
--- /dev/null
+++ b/parquet/pqarrow/file_reader_bench_test.go
@@ -0,0 +1,104 @@
+// 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 pqarrow_test
+
+import (
+ "bytes"
+ "context"
+ "fmt"
+ "testing"
+
+ "github.com/apache/arrow-go/v18/arrow"
+ "github.com/apache/arrow-go/v18/arrow/array"
+ "github.com/apache/arrow-go/v18/arrow/memory"
+ "github.com/apache/arrow-go/v18/parquet/file"
+ "github.com/apache/arrow-go/v18/parquet/pqarrow"
+)
+
+func BenchmarkReadTableSerial(b *testing.B) {
+ const (
+ rows = 1024
+ rowGroupSize = 64
+ )
+
+ for _, numColumns := range []int{1, 8, 32, 128} {
+ b.Run(fmt.Sprintf("columns=%d", numColumns), func(b *testing.B)
{
+ mem := memory.DefaultAllocator
+ tbl := makeWideInt32Table(mem, numColumns, rows)
+ defer tbl.Release()
+
+ var buf bytes.Buffer
+ if err := pqarrow.WriteTable(tbl, &buf, rowGroupSize,
nil, pqarrow.DefaultWriterProps()); err != nil {
+ b.Fatal(err)
+ }
+ parquetData := buf.Bytes()
+
+ b.ReportAllocs()
+ b.SetBytes(int64(len(parquetData)))
+ b.ResetTimer()
+ for range b.N {
+ pf, err :=
file.NewParquetReader(bytes.NewReader(parquetData))
+ if err != nil {
+ b.Fatal(err)
+ }
+
+ reader, err := pqarrow.NewFileReader(pf,
pqarrow.ArrowReadProperties{Parallel: false}, mem)
+ if err != nil {
+ _ = pf.Close()
+ b.Fatal(err)
+ }
+
+ out, err :=
reader.ReadTable(context.Background())
+ if err != nil {
+ _ = pf.Close()
+ b.Fatal(err)
+ }
+ out.Release()
+ if err := pf.Close(); err != nil {
+ b.Fatal(err)
+ }
+ }
+ })
+ }
+}
+
+func makeWideInt32Table(mem memory.Allocator, numColumns, numRows int)
arrow.Table {
+ values := make([]int32, numRows)
+ for i := range values {
+ values[i] = int32(i)
+ }
+
+ fields := make([]arrow.Field, numColumns)
+ columns := make([]arrow.Column, numColumns)
+ for i := range columns {
+ fields[i] = arrow.Field{Name: fmt.Sprintf("column_%d", i),
Type: arrow.PrimitiveTypes.Int32}
+
+ builder := array.NewInt32Builder(mem)
+ builder.AppendValues(values, nil)
+ arr := builder.NewInt32Array()
+ builder.Release()
+
+ columns[i] = arrow.NewColumnFromArr(fields[i], arr)
+ arr.Release()
+ }
+
+ table := array.NewTable(arrow.NewSchema(fields, nil), columns,
int64(numRows))
+ for i := range columns {
+ columns[i].Release()
+ }
+ return table
+}
diff --git a/parquet/pqarrow/file_reader_extension_test.go
b/parquet/pqarrow/file_reader_extension_test.go
index 76429e8a..4adbd6ed 100644
--- a/parquet/pqarrow/file_reader_extension_test.go
+++ b/parquet/pqarrow/file_reader_extension_test.go
@@ -17,12 +17,17 @@
package pqarrow
import (
+ "bytes"
+ "context"
"reflect"
+ "sync"
"testing"
"github.com/apache/arrow-go/v18/arrow"
"github.com/apache/arrow-go/v18/arrow/array"
"github.com/apache/arrow-go/v18/arrow/memory"
+ "github.com/apache/arrow-go/v18/parquet"
+ "github.com/apache/arrow-go/v18/parquet/file"
"github.com/stretchr/testify/require"
)
@@ -43,6 +48,33 @@ type stableExtensionArray struct {
array.ExtensionArrayBase
}
+type panickingExtensionType struct {
+ arrow.ExtensionBase
+ beforePanic func()
+}
+
+func (t *panickingExtensionType) ArrayType() reflect.Type {
+ if t.beforePanic != nil {
+ t.beforePanic()
+ }
+ panic("malformed extension array type")
+}
+
+func (*panickingExtensionType) ExtensionName() string { return
"test.panicking" }
+
+func (*panickingExtensionType) ExtensionEquals(other arrow.ExtensionType) bool
{
+ _, ok := other.(*panickingExtensionType)
+ return ok
+}
+
+func (*panickingExtensionType) Serialize() string { return "" }
+
+func (*panickingExtensionType) Deserialize(arrow.DataType, string)
(arrow.ExtensionType, error) {
+ return &panickingExtensionType{
+ ExtensionBase: arrow.ExtensionBase{Storage:
arrow.PrimitiveTypes.Int32},
+ }, nil
+}
+
func (*stableExtensionType) StorageType() arrow.DataType { return
arrow.PrimitiveTypes.Int32 }
func (*stableExtensionType) ArrayType() reflect.Type {
@@ -156,3 +188,161 @@ func TestExtensionReaderBuildArrayReleasesChunks(t
*testing.T) {
require.NoError(t, err)
out.Release()
}
+
+func TestReadRowGroupsSerialRecoversFromPanic(t *testing.T) {
+ mem := memory.NewCheckedAllocator(memory.NewGoAllocator())
+ defer mem.AssertSize(t, 0)
+
+ schema := arrow.NewSchema([]arrow.Field{{Name: "value", Type:
arrow.PrimitiveTypes.Int32}}, nil)
+ builder := array.NewInt32Builder(mem)
+ builder.Append(1)
+ values := builder.NewInt32Array()
+ builder.Release()
+ defer values.Release()
+
+ record := array.NewRecordBatch(schema, []arrow.Array{values}, 1)
+ defer record.Release()
+
+ var buf bytes.Buffer
+ writer, err := NewFileWriter(schema, &buf, nil, DefaultWriterProps())
+ require.NoError(t, err)
+ require.NoError(t, writer.Write(record))
+ require.NoError(t, writer.Close())
+
+ parquetReader, err :=
file.NewParquetReader(bytes.NewReader(buf.Bytes()),
+ file.WithReadProps(parquet.NewReaderProperties(mem)))
+ require.NoError(t, err)
+ defer parquetReader.Close()
+
+ reader, err := NewFileReader(parquetReader,
ArrowReadProperties{Parallel: false}, mem)
+ require.NoError(t, err)
+ reader.Manifest.Fields[0].Field.Type = &panickingExtensionType{
+ ExtensionBase: arrow.ExtensionBase{Storage:
arrow.PrimitiveTypes.Int32},
+ }
+
+ var table arrow.Table
+ var readErr error
+ require.NotPanics(t, func() {
+ table, readErr = reader.ReadRowGroups(context.Background(),
[]int{0}, []int{0})
+ })
+ require.Nil(t, table)
+ require.ErrorContains(t, readErr, "panic while reading")
+ require.ErrorContains(t, readErr, "malformed extension array type")
+}
+
+func TestReadRowGroupsRecoversFromMultiplePanics(t *testing.T) {
+ for _, parallel := range []bool{false, true} {
+ name := "serial"
+ if parallel {
+ name = "parallel"
+ }
+ t.Run(name, func(t *testing.T) {
+ mem :=
memory.NewCheckedAllocator(memory.DefaultAllocator)
+ defer mem.AssertSize(t, 0)
+
+ schema := arrow.NewSchema([]arrow.Field{
+ {Name: "valid", Type:
arrow.PrimitiveTypes.Int32},
+ {Name: "invalid1", Type:
arrow.PrimitiveTypes.Int32},
+ {Name: "invalid2", Type:
arrow.PrimitiveTypes.Int32},
+ }, nil)
+ builder := array.NewInt32Builder(mem)
+ builder.Append(1)
+ values := builder.NewInt32Array()
+ builder.Release()
+ defer values.Release()
+ record := array.NewRecordBatch(schema,
[]arrow.Array{values, values, values}, 1)
+ defer record.Release()
+
+ var buf bytes.Buffer
+ writer, err := NewFileWriter(schema, &buf, nil,
DefaultWriterProps())
+ require.NoError(t, err)
+ require.NoError(t, writer.Write(record))
+ require.NoError(t, writer.Close())
+
+ parquetReader, err :=
file.NewParquetReader(bytes.NewReader(buf.Bytes()),
+
file.WithReadProps(parquet.NewReaderProperties(mem)))
+ require.NoError(t, err)
+ defer parquetReader.Close()
+ reader, err := NewFileReader(parquetReader,
ArrowReadProperties{Parallel: parallel}, mem)
+ require.NoError(t, err)
+
+ var ready sync.WaitGroup
+ ready.Add(2)
+ typ := &panickingExtensionType{
+ ExtensionBase: arrow.ExtensionBase{Storage:
arrow.PrimitiveTypes.Int32},
+ }
+ if parallel {
+ typ.beforePanic = func() {
+ ready.Done()
+ ready.Wait()
+ }
+ }
+ reader.Manifest.Fields[1].Field.Type = typ
+ reader.Manifest.Fields[2].Field.Type = typ
+
+ var table arrow.Table
+ var readErr error
+ require.NotPanics(t, func() {
+ table, readErr =
reader.ReadTable(context.Background())
+ })
+ require.Nil(t, table)
+ require.ErrorContains(t, readErr, "panic while reading")
+ require.ErrorContains(t, readErr, "malformed extension
array type")
+ })
+ }
+}
+
+type cancelingExtensionType struct {
+ stableExtensionType
+ cancel context.CancelFunc
+}
+
+func (t *cancelingExtensionType) ArrayType() reflect.Type {
+ t.cancel()
+ return reflect.TypeFor[stableExtensionArray]()
+}
+
+func TestReadRowGroupsCanceledDuringFinalColumn(t *testing.T) {
+ mem := memory.NewCheckedAllocator(memory.DefaultAllocator)
+ defer mem.AssertSize(t, 0)
+
+ schema := arrow.NewSchema([]arrow.Field{{Name: "value", Type:
arrow.PrimitiveTypes.Int32}}, nil)
+ builder := array.NewInt32Builder(mem)
+ builder.Append(1)
+ values := builder.NewInt32Array()
+ builder.Release()
+ defer values.Release()
+ record := array.NewRecordBatch(schema, []arrow.Array{values}, 1)
+ defer record.Release()
+
+ var buf bytes.Buffer
+ writer, err := NewFileWriter(schema, &buf, nil, DefaultWriterProps())
+ require.NoError(t, err)
+ require.NoError(t, writer.Write(record))
+ require.NoError(t, writer.Close())
+
+ for _, parallel := range []bool{false, true} {
+ name := "serial"
+ if parallel {
+ name = "parallel"
+ }
+ t.Run(name, func(t *testing.T) {
+ parquetReader, err :=
file.NewParquetReader(bytes.NewReader(buf.Bytes()),
+
file.WithReadProps(parquet.NewReaderProperties(mem)))
+ require.NoError(t, err)
+ defer parquetReader.Close()
+ reader, err := NewFileReader(parquetReader,
ArrowReadProperties{Parallel: parallel}, mem)
+ require.NoError(t, err)
+ ctx, cancel := context.WithCancel(context.Background())
+ defer cancel()
+ reader.Manifest.Fields[0].Field.Type =
&cancelingExtensionType{cancel: cancel}
+
+ table, err := reader.ReadTable(ctx)
+ if table != nil {
+ defer table.Release()
+ }
+ require.ErrorIs(t, err, context.Canceled)
+ require.Nil(t, table)
+ })
+ }
+}