From f6b03061baedba6bf2895813580add00920b595c Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Thu, 27 Aug 2026 15:09:12 +0200 Subject: [PATCH 1/2] perf(parquet/pqarrow): use synchronous path for serial reads --- parquet/pqarrow/file_reader.go | 54 ++++++++++- parquet/pqarrow/file_reader_bench_test.go | 104 ++++++++++++++++++++++ 2 files changed, 155 insertions(+), 3 deletions(-) create mode 100644 parquet/pqarrow/file_reader_bench_test.go diff --git a/parquet/pqarrow/file_reader.go b/parquet/pqarrow/file_reader.go index ae659d220..88f983622 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 { @@ -370,6 +385,39 @@ 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() + } + + 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 diff --git a/parquet/pqarrow/file_reader_bench_test.go b/parquet/pqarrow/file_reader_bench_test.go new file mode 100644 index 000000000..84578c0ce --- /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 +} From 871a6445db210f635e018c18bcba4aca24a6b072 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Sat, 29 Aug 2026 11:35:02 +0200 Subject: [PATCH 2/2] fix(parquet/pqarrow): recover serial read panics --- parquet/pqarrow/file_reader.go | 13 +++- parquet/pqarrow/file_reader_extension_test.go | 68 +++++++++++++++++++ 2 files changed, 79 insertions(+), 2 deletions(-) diff --git a/parquet/pqarrow/file_reader.go b/parquet/pqarrow/file_reader.go index 88f983622..6243ac495 100644 --- a/parquet/pqarrow/file_reader.go +++ b/parquet/pqarrow/file_reader.go @@ -320,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 ( @@ -399,7 +408,7 @@ func (fr *FileReader) ReadRowGroups(ctx context.Context, indices, rowGroups []in return nil, err } - data, err := fr.ReadColumn(rowGroups, rdr) + data, err := fr.readColumn(rowGroups, rdr) if err != nil { if data != nil { data.Release() @@ -451,7 +460,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} diff --git a/parquet/pqarrow/file_reader_extension_test.go b/parquet/pqarrow/file_reader_extension_test.go index 76429e8a5..a00fc9e70 100644 --- a/parquet/pqarrow/file_reader_extension_test.go +++ b/parquet/pqarrow/file_reader_extension_test.go @@ -17,12 +17,16 @@ package pqarrow import ( + "bytes" + "context" "reflect" "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 +47,29 @@ type stableExtensionArray struct { array.ExtensionArrayBase } +type panickingExtensionType struct { + arrow.ExtensionBase +} + +func (*panickingExtensionType) ArrayType() reflect.Type { + 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 +183,44 @@ 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") +}