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 e0969720 fix(parquet/pqarrow): release partially constructed field 
readers (#1061)
e0969720 is described below

commit e0969720e786504e371085fd1b603f075bd8e669
Author: Minh Vu <[email protected]>
AuthorDate: Wed Aug 5 19:45:35 2026 +0200

    fix(parquet/pqarrow): release partially constructed field readers (#1061)
    
    ### Rationale for this change
    
    `GetFieldReaders` constructs readers concurrently. If one fails after
    another succeeds, the successful reader remains in a private slice and
    is never returned, so the caller cannot release it.
    
    ### What changes are included in this PR?
    
    When construction fails, release every non-nil reader that completed
    before returning the original error.
    
    ### Are these changes tested?
    
    Yes. The regression test injects a read failure at the second column
    chunk after the first reader has been constructed. Without the cleanup,
    the checked allocator reports the first page reader's retained buffer;
    with the fix, allocator usage returns to zero.
    
    `go test ./parquet/pqarrow -run
    TestGetFieldReadersReleasesPartialReadersOnError`
    
    ### Are there any user-facing changes?
    
    No API changes. Failed multi-column reader construction now cleans up
    partial results.
---
 parquet/pqarrow/file_reader.go      |  5 ++++
 parquet/pqarrow/file_reader_test.go | 60 +++++++++++++++++++++++++++++++++++++
 2 files changed, 65 insertions(+)

diff --git a/parquet/pqarrow/file_reader.go b/parquet/pqarrow/file_reader.go
index ed6a2998..96af9f77 100644
--- a/parquet/pqarrow/file_reader.go
+++ b/parquet/pqarrow/file_reader.go
@@ -262,6 +262,11 @@ func (fr *FileReader) GetFieldReaders(ctx context.Context, 
colIndices, rowGroups
                })
        }
        if err = g.Wait(); err != nil {
+               for _, rdr := range out {
+                       if rdr != nil {
+                               rdr.Release()
+                       }
+               }
                return nil, nil, err
        }
 
diff --git a/parquet/pqarrow/file_reader_test.go 
b/parquet/pqarrow/file_reader_test.go
index aeb32086..45e0a4f3 100644
--- a/parquet/pqarrow/file_reader_test.go
+++ b/parquet/pqarrow/file_reader_test.go
@@ -19,6 +19,7 @@ package pqarrow_test
 import (
        "bytes"
        "context"
+       "errors"
        "fmt"
        "io"
        "os"
@@ -191,6 +192,65 @@ func TestArrowReaderCanceledContext(t *testing.T) {
        require.ErrorIs(t, err, context.Canceled)
 }
 
+type failingReaderAt struct {
+       *bytes.Reader
+       failOffset int64
+       err        error
+}
+
+func (r *failingReaderAt) ReadAt(p []byte, off int64) (int, error) {
+       if off == r.failOffset {
+               return 0, r.err
+       }
+       return r.Reader.ReadAt(p, off)
+}
+
+func TestGetFieldReadersReleasesPartialReadersOnError(t *testing.T) {
+       schema := arrow.NewSchema([]arrow.Field{
+               {Name: "first", Type: arrow.PrimitiveTypes.Int32},
+               {Name: "second", Type: arrow.PrimitiveTypes.Int32},
+       }, nil)
+       record, _, err := array.RecordFromJSON(memory.DefaultAllocator, schema,
+               strings.NewReader(`[{"first": 1, "second": 2}]`))
+       require.NoError(t, err)
+       defer record.Release()
+
+       var buf bytes.Buffer
+       writer, err := pqarrow.NewFileWriter(schema, &buf, nil, 
pqarrow.DefaultWriterProps())
+       require.NoError(t, err)
+       require.NoError(t, writer.Write(record))
+       require.NoError(t, writer.Close())
+
+       readErr := errors.New("read second column")
+       source := &failingReaderAt{
+               Reader:     bytes.NewReader(buf.Bytes()),
+               failOffset: -1,
+               err:        readErr,
+       }
+       mem := memory.NewCheckedAllocator(memory.DefaultAllocator)
+       defer mem.AssertSize(t, 0)
+
+       parquetReader, err := file.NewParquetReader(source,
+               file.WithReadProps(parquet.NewReaderProperties(mem)))
+       require.NoError(t, err)
+       defer parquetReader.Close()
+
+       column, err := parquetReader.MetaData().RowGroup(0).ColumnChunk(1)
+       require.NoError(t, err)
+       source.failOffset = column.DataPageOffset()
+       if column.HasDictionaryPage() && column.DictionaryPageOffset() > 0 {
+               source.failOffset = column.DictionaryPageOffset()
+       }
+
+       reader, err := pqarrow.NewFileReader(parquetReader, 
pqarrow.ArrowReadProperties{}, mem)
+       require.NoError(t, err)
+       readers, resultSchema, err := 
reader.GetFieldReaders(context.Background(), []int{0, 1}, []int{0})
+       require.ErrorIs(t, err, readErr)
+       require.Nil(t, readers)
+       require.Nil(t, resultSchema)
+       require.Zero(t, mem.CurrentAlloc())
+}
+
 func TestRecordReaderParallel(t *testing.T) {
        mem := memory.NewCheckedAllocator(memory.DefaultAllocator)
        defer mem.AssertSize(t, 0)

Reply via email to