From 4a79daa6d6b6335339582d88e6f59d3b7255ddbb Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Fri, 28 Aug 2026 12:56:04 +0200 Subject: [PATCH 01/10] perf(table): load position deletes lazily per task Signed-off-by: Minh Vu --- table/arrow_scanner.go | 243 +++++++++++---- table/arrow_scanner_lazy_delete_bench_test.go | 165 +++++++++++ ...row_scanner_lazy_delete_regression_test.go | 278 ++++++++++++++++++ 3 files changed, 627 insertions(+), 59 deletions(-) create mode 100644 table/arrow_scanner_lazy_delete_bench_test.go create mode 100644 table/arrow_scanner_lazy_delete_regression_test.go diff --git a/table/arrow_scanner.go b/table/arrow_scanner.go index 73a3cf8ff..16d8ba505 100644 --- a/table/arrow_scanner.go +++ b/table/arrow_scanner.go @@ -135,6 +135,113 @@ func readAllDeleteFiles(ctx context.Context, fs iceio.IO, tasks []FileScanTask, return deletesPerFile, nil } +// lazyPositionDeleteLoader indexes positional-delete metadata for a scan, but +// waits to open each delete file until a worker reaches a task that references +// it. A delete file can apply to more than one data file, so the cache keeps +// the complete grouped result for the delete file rather than caching only one +// task's positions. +// +// The grouped Arrow chunks are owned by the loader until release. The iterator +// calls release after all workers have stopped, which keeps shared chunks alive +// while multiple tasks use them and also covers early iterator termination. +type lazyPositionDeleteLoader struct { + fs iceio.IO + files map[string]*lazyPositionDeleteFile + + releaseOnce sync.Once +} + +type lazyPositionDeleteFile struct { + dataFile iceberg.DataFile + + once sync.Once + deletes map[string]*arrow.Chunked + err error +} + +func newLazyPositionDeleteLoader(fs iceio.IO, tasks []FileScanTask) *lazyPositionDeleteLoader { + loader := &lazyPositionDeleteLoader{ + fs: fs, + files: make(map[string]*lazyPositionDeleteFile), + } + + for _, task := range tasks { + for _, deleteFile := range task.DeleteFiles { + if deleteFile.ContentType() != iceberg.EntryContentPosDeletes { + continue + } + + path := deleteFile.FilePath() + if _, ok := loader.files[path]; !ok { + loader.files[path] = &lazyPositionDeleteFile{dataFile: deleteFile} + } + } + } + + return loader +} + +func (l *lazyPositionDeleteLoader) load(ctx context.Context, task FileScanTask) (positionDeletes, error) { + if len(task.DeleteFiles) == 0 { + return nil, nil + } + + targetPath := task.File.FilePath() + deletes := make(positionDeletes, 0, len(task.DeleteFiles)) + seen := make(map[string]struct{}, len(task.DeleteFiles)) + for _, deleteFile := range task.DeleteFiles { + if deleteFile.ContentType() != iceberg.EntryContentPosDeletes { + continue + } + + path := deleteFile.FilePath() + if _, ok := seen[path]; ok { + continue + } + seen[path] = struct{}{} + + cached, ok := l.files[path] + if !ok { + // The loader is normally built from the same task slice supplied to + // this method. Keep this guard so a malformed caller cannot panic a + // scan if it changes a task after loader construction. + continue + } + + cached.once.Do(func() { + cached.deletes, cached.err = readDeletes(ctx, l.fs, cached.dataFile) + if cached.err != nil { + // readDeletes currently returns nil on errors. Release defensively + // in case a future reader returns partial Arrow ownership. + releasePosDeletes(cached.deletes) + cached.deletes = nil + } + }) + if cached.err != nil { + return nil, cached.err + } + + if chunk := cached.deletes[targetPath]; chunk != nil { + deletes = append(deletes, chunk) + } + } + + return deletes, nil +} + +func (l *lazyPositionDeleteLoader) release() { + if l == nil { + return + } + + l.releaseOnce.Do(func() { + for _, cached := range l.files { + releasePosDeletes(cached.deletes) + cached.deletes = nil + } + }) +} + // perFileDVBitmaps maps each data-file path to the deletion-vector bitmap // that applies to it. Kept separate from perFilePosDeletes so the row-filter // pipeline can use compute.Filter on a Boolean mask built from Contains() @@ -1702,6 +1809,10 @@ func (as *arrowScan) producePosDeletesFromTask(ctx context.Context, task tblutil } func createIterator(ctx context.Context, numWorkers uint, records <-chan enumeratedRecord, deletesPerFile perFilePosDeletes, cancel context.CancelCauseFunc, rowLimit int64) iter.Seq2[arrow.RecordBatch, error] { + return createIteratorWithCleanup(ctx, numWorkers, records, deletesPerFile, cancel, rowLimit, nil) +} + +func createIteratorWithCleanup(ctx context.Context, numWorkers uint, records <-chan enumeratedRecord, deletesPerFile perFilePosDeletes, cancel context.CancelCauseFunc, rowLimit int64, cleanup func()) iter.Seq2[arrow.RecordBatch, error] { isBeforeAny := func(batch enumeratedRecord) bool { return batch.Task.Index < 0 } @@ -1745,6 +1856,9 @@ func createIterator(ctx context.Context, numWorkers uint, records <-chan enumera } releasePerFilePosDeletes(deletesPerFile) + if cleanup != nil { + cleanup() + } }() defer cancel(nil) @@ -1799,59 +1913,81 @@ func createIterator(ctx context.Context, numWorkers uint, records <-chan enumera } } -func (as *arrowScan) recordBatchesFromTasksAndDeletes(ctx context.Context, tasks []FileScanTask, deletesPerFile perFilePosDeletes, dvBitmaps perFileDVBitmaps, eqDeleteSets map[int][]*equalityDeleteSet, invariants *arrowScanInvariants) iter.Seq2[arrow.RecordBatch, error] { - extSet := substrait.NewExtensionSet() - - ctx, cancel := context.WithCancelCause(exprs.WithExtensionIDSet(ctx, extSet)) - taskChan := make(chan tblutils.Enumerated[FileScanTask], len(tasks)) - - // numWorkers := 1 - numWorkers := min(as.concurrency, len(tasks)) - records := make(chan enumeratedRecord, numWorkers) +func (as *arrowScan) recordBatchesFromTasksAndDeletes(ctx context.Context, tasks []FileScanTask, positionDeleteLoader *lazyPositionDeleteLoader, dvBitmaps perFileDVBitmaps, eqDeleteSets map[int][]*equalityDeleteSet, invariants *arrowScanInvariants) iter.Seq2[arrow.RecordBatch, error] { + return func(yield func(arrow.RecordBatch, error) bool) { + extSet := substrait.NewExtensionSet() + scanCtx, cancel := context.WithCancelCause(exprs.WithExtensionIDSet(ctx, extSet)) + numWorkers := min(as.concurrency, len(tasks)) + taskChan := make(chan tblutils.Enumerated[FileScanTask], len(tasks)) + records := make(chan enumeratedRecord, numWorkers) + + var wg sync.WaitGroup + wg.Add(numWorkers) + for range numWorkers { + go func() { + defer wg.Done() + for { + select { + case <-scanCtx.Done(): + return + case task, ok := <-taskChan: + if !ok { + return + } + if scanCtx.Err() != nil { + return + } + + filePath := task.Value.File.FilePath() + var positionalDeletes positionDeletes + if positionDeleteLoader != nil { + var err error + positionalDeletes, err = positionDeleteLoader.load(scanCtx, task.Value) + if err != nil { + records <- enumeratedRecord{Task: task, Err: err} + cancel(err) + + return + } + } + + if err := as.recordsFromTask(scanCtx, task, records, + positionalDeletes, + dvBitmaps[filePath], + eqDeleteSets[task.Index], + invariants); err != nil { + cancel(err) + + return + } + } + } + }() + } - var wg sync.WaitGroup - wg.Add(numWorkers) - for range numWorkers { go func() { - defer wg.Done() - for { + for i, t := range tasks { select { - case <-ctx.Done(): + case <-scanCtx.Done(): return - case task, ok := <-taskChan: - if !ok { - return - } - - filePath := task.Value.File.FilePath() - if err := as.recordsFromTask(ctx, task, records, - deletesPerFile[filePath], - dvBitmaps[filePath], - eqDeleteSets[task.Index], - invariants); err != nil { - cancel(err) - - return - } + case taskChan <- tblutils.Enumerated[FileScanTask]{ + Value: t, Index: i, Last: i == len(tasks)-1, + }: } } + close(taskChan) + + wg.Wait() + close(records) }() - } - go func() { - for i, t := range tasks { - taskChan <- tblutils.Enumerated[FileScanTask]{ - Value: t, Index: i, Last: i == len(tasks)-1, - } + var cleanup func() + if positionDeleteLoader != nil { + cleanup = positionDeleteLoader.release } - close(taskChan) - - wg.Wait() - close(records) - }() - - return createIterator(ctx, uint(numWorkers), records, deletesPerFile, - cancel, as.rowLimit) + createIteratorWithCleanup(scanCtx, uint(numWorkers), records, nil, + cancel, as.rowLimit, cleanup)(yield) + } } func (as *arrowScan) GetRecords(ctx context.Context, tasks []FileScanTask) (*arrow.Schema, iter.Seq2[arrow.RecordBatch, error], error) { @@ -1884,15 +2020,6 @@ func (as *arrowScan) GetRecords(ctx context.Context, tasks []FileScanTask) (*arr return nil, nil, err } - deletesPerFile, err := readAllDeleteFiles(ctx, as.fs, tasks, as.concurrency) - if err != nil { - // readAllDeleteFiles can return a partially-populated map alongside - // the error if some goroutines completed before the failure. - releasePerFilePosDeletes(deletesPerFile) - - return nil, nil, err - } - // DV bitmaps stay in their native form rather than being materialized // into int64 positions and merged with the Parquet pos-delete map. // filterByDeletionVector applies the bitmap to each batch via a Boolean @@ -1900,20 +2027,18 @@ func (as *arrowScan) GetRecords(ctx context.Context, tasks []FileScanTask) (*arr // no intermediate position set. dvBitmaps, err := readAllDeletionVectors(ctx, as.fs, tasks, as.concurrency) if err != nil { - releasePerFilePosDeletes(deletesPerFile) - return nil, nil, err } eqDeleteSets, err := readAllEqualityDeleteFiles(ctx, as.fs, invariants.tableSchema, invariants.nameMapping, tasks, as.concurrency) if err != nil { - // Positional deletes were fully loaded; release them before aborting. - releasePerFilePosDeletes(deletesPerFile) - return nil, nil, err } addEqualityDeleteFieldIDs(invariants, eqDeleteSets) - return resultSchema, as.recordBatchesFromTasksAndDeletes(ctx, tasks, deletesPerFile, dvBitmaps, eqDeleteSets, invariants), nil + positionDeleteLoader := newLazyPositionDeleteLoader(as.fs, tasks) + + return resultSchema, as.recordBatchesFromTasksAndDeletes(ctx, tasks, + positionDeleteLoader, dvBitmaps, eqDeleteSets, invariants), nil } diff --git a/table/arrow_scanner_lazy_delete_bench_test.go b/table/arrow_scanner_lazy_delete_bench_test.go new file mode 100644 index 000000000..da3a22f92 --- /dev/null +++ b/table/arrow_scanner_lazy_delete_bench_test.go @@ -0,0 +1,165 @@ +// 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 table + +import ( + "fmt" + "strings" + "testing" + + "github.com/apache/arrow-go/v18/arrow" + "github.com/apache/arrow-go/v18/arrow/array" + "github.com/apache/arrow-go/v18/parquet" + "github.com/apache/arrow-go/v18/parquet/pqarrow" + "github.com/apache/iceberg-go" + iceio "github.com/apache/iceberg-go/io" +) + +const ( + lazyPositionDeleteBenchmarkTaskCount = 10_000 + lazyPositionDeleteBenchmarkFileCount = 1_000 + lazyPositionDeleteBenchmarkTasksPerFile = lazyPositionDeleteBenchmarkTaskCount / lazyPositionDeleteBenchmarkFileCount + lazyPositionDeleteBenchmarkDataFileBytes = 128 +) + +type lazyPositionDeleteBenchmarkFixture struct { + fs *iceio.MemFS + tasks []FileScanTask +} + +func newLazyPositionDeleteBenchmarkFixture(b *testing.B) lazyPositionDeleteBenchmarkFixture { + b.Helper() + + fs := iceio.NewMemFS() + deleteFiles := make([]iceberg.DataFile, lazyPositionDeleteBenchmarkFileCount) + for deleteIndex := range deleteFiles { + deletePath := fmt.Sprintf("mem://benchmark/deletes/delete-%04d.parquet", deleteIndex) + var content strings.Builder + content.WriteByte('[') + for taskOffset := range lazyPositionDeleteBenchmarkTasksPerFile { + if taskOffset > 0 { + content.WriteByte(',') + } + taskIndex := deleteIndex + taskOffset*lazyPositionDeleteBenchmarkFileCount + fmt.Fprintf(&content, + `{"file_path":"mem://benchmark/data/data-%05d.parquet","pos":0}`, + taskIndex) + } + content.WriteByte(']') + + benchmarkWritePosDeleteParquet(b, fs, deletePath, content.String()) + builder, err := iceberg.NewDataFileBuilder( + *iceberg.UnpartitionedSpec, iceberg.EntryContentPosDeletes, + deletePath, iceberg.ParquetFile, nil, nil, nil, + lazyPositionDeleteBenchmarkTasksPerFile, lazyPositionDeleteBenchmarkDataFileBytes) + if err != nil { + b.Fatal(err) + } + deleteFiles[deleteIndex] = builder.Build() + } + + tasks := make([]FileScanTask, lazyPositionDeleteBenchmarkTaskCount) + for taskIndex := range tasks { + dataPath := fmt.Sprintf("mem://benchmark/data/data-%05d.parquet", taskIndex) + builder, err := iceberg.NewDataFileBuilder( + *iceberg.UnpartitionedSpec, iceberg.EntryContentData, + dataPath, iceberg.ParquetFile, nil, nil, nil, 1, lazyPositionDeleteBenchmarkDataFileBytes) + if err != nil { + b.Fatal(err) + } + tasks[taskIndex] = FileScanTask{ + File: builder.Build(), + DeleteFiles: []iceberg.DataFile{deleteFiles[taskIndex%lazyPositionDeleteBenchmarkFileCount]}, + } + } + + return lazyPositionDeleteBenchmarkFixture{fs: fs, tasks: tasks} +} + +func benchmarkWritePosDeleteParquet(b *testing.B, fs *iceio.MemFS, path, content string) { + b.Helper() + + record := mustLoadRecordBatchFromJSON(PositionalDeleteArrowSchema, content) + defer record.Release() + tbl := array.NewTableFromRecords(PositionalDeleteArrowSchema, []arrow.RecordBatch{record}) + defer tbl.Release() + + file, err := fs.Create(path) + if err != nil { + b.Fatal(err) + } + if err := pqarrow.WriteTable(tbl, file, record.NumRows(), + parquet.NewWriterProperties(parquet.WithStats(true)), + pqarrow.DefaultWriterProps()); err != nil { + b.Fatal(err) + } + if err := file.Close(); err != nil { + b.Fatal(err) + } +} + +func BenchmarkLazyPositionDeleteLoading(b *testing.B) { + fixture := newLazyPositionDeleteBenchmarkFixture(b) + + b.Run("eager_all_delete_files", func(b *testing.B) { + b.ReportAllocs() + b.ResetTimer() + for b.Loop() { + deletes, err := readAllDeleteFiles(b.Context(), fixture.fs, fixture.tasks, 16) + if err != nil { + b.Fatal(err) + } + releasePerFilePosDeletes(deletes) + } + b.ReportMetric(lazyPositionDeleteBenchmarkFileCount, "delete_files/op") + b.ReportMetric(lazyPositionDeleteBenchmarkFileCount, "delete_files_before_first/op") + }) + + b.Run("lazy_unread_iterator", func(b *testing.B) { + b.ReportAllocs() + b.ResetTimer() + for b.Loop() { + loader := newLazyPositionDeleteLoader(fixture.fs, fixture.tasks) + if len(loader.files) != lazyPositionDeleteBenchmarkFileCount { + b.Fatalf("expected %d cached delete files, got %d", + lazyPositionDeleteBenchmarkFileCount, len(loader.files)) + } + loader.release() + } + b.ReportMetric(0, "delete_files/op") + b.ReportMetric(0, "delete_files_before_first/op") + }) + + b.Run("lazy_first_task", func(b *testing.B) { + b.ReportAllocs() + b.ResetTimer() + for b.Loop() { + loader := newLazyPositionDeleteLoader(fixture.fs, fixture.tasks) + deletes, err := loader.load(b.Context(), fixture.tasks[0]) + if err != nil { + b.Fatal(err) + } + if len(deletes) != 1 { + b.Fatalf("expected one delete chunk for the first task, got %d", len(deletes)) + } + loader.release() + } + b.ReportMetric(1, "delete_files/op") + b.ReportMetric(1, "delete_files_before_first/op") + }) +} diff --git a/table/arrow_scanner_lazy_delete_regression_test.go b/table/arrow_scanner_lazy_delete_regression_test.go new file mode 100644 index 000000000..04b69ed2e --- /dev/null +++ b/table/arrow_scanner_lazy_delete_regression_test.go @@ -0,0 +1,278 @@ +// 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 table + +import ( + "context" + "sync" + "sync/atomic" + "testing" + + "github.com/apache/arrow-go/v18/arrow" + "github.com/apache/arrow-go/v18/arrow/compute" + "github.com/apache/arrow-go/v18/arrow/memory" + "github.com/apache/iceberg-go" + iceio "github.com/apache/iceberg-go/io" + "github.com/apache/iceberg-go/table/internal" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type countingOpenMemFS struct { + *iceio.MemFS + opens atomic.Int64 +} + +func (f *countingOpenMemFS) Open(name string) (iceio.File, error) { + f.opens.Add(1) + + return f.MemFS.Open(name) +} + +func newLazyDataFile(t *testing.T, path string) iceberg.DataFile { + t.Helper() + + builder, err := iceberg.NewDataFileBuilder( + *iceberg.UnpartitionedSpec, iceberg.EntryContentData, + path, iceberg.ParquetFile, nil, nil, nil, 1, 128) + require.NoError(t, err) + + return builder.Build() +} + +func TestLazyPositionDeleteLoaderDefersReadsAndSharesResults(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + ctx := compute.WithAllocator(t.Context(), mem) + + const ( + deletePath = "mem://bucket/deletes/shared.parquet" + dataPathA = "mem://bucket/data/a.parquet" + dataPathB = "mem://bucket/data/b.parquet" + ) + memFS := &countingOpenMemFS{MemFS: iceio.NewMemFS()} + writePosDeleteParquetToMemFS(t, memFS.MemFS, deletePath, `[ + {"file_path": "`+dataPathA+`", "pos": 1}, + {"file_path": "`+dataPathB+`", "pos": 3} + ]`) + + deleteFile := newPosDeleteFile(t, deletePath, 2, 128) + tasks := []FileScanTask{ + { + File: newLazyDataFile(t, dataPathA), + DeleteFiles: []iceberg.DataFile{deleteFile, deleteFile}, + }, + { + File: newLazyDataFile(t, dataPathB), + DeleteFiles: []iceberg.DataFile{deleteFile}, + }, + } + loader := newLazyPositionDeleteLoader(memFS, tasks) + + assert.Zero(t, memFS.opens.Load(), "constructing the scan loader must not open delete files") + + gotA, err := loader.load(ctx, tasks[0]) + require.NoError(t, err) + require.Len(t, gotA, 1, "duplicate delete references must be read once per task") + assert.Equal(t, []int64{1}, int64Values(gotA[0])) + assert.Equal(t, int64(1), memFS.opens.Load()) + + gotB, err := loader.load(ctx, tasks[1]) + require.NoError(t, err) + require.Len(t, gotB, 1) + assert.Equal(t, []int64{3}, int64Values(gotB[0])) + assert.Equal(t, int64(1), memFS.opens.Load(), "shared delete files must use one read") + + loader.release() +} + +func TestLazyPositionDeleteLoaderSingleflightsConcurrentLoads(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + ctx := compute.WithAllocator(t.Context(), mem) + + const ( + deletePath = "mem://bucket/deletes/concurrent.parquet" + dataPath = "mem://bucket/data/a.parquet" + ) + fs := &countingOpenMemFS{MemFS: iceio.NewMemFS()} + writePosDeleteParquetToMemFS(t, fs.MemFS, deletePath, + `[{"file_path":"`+dataPath+`","pos":7}]`) + deleteFile := newPosDeleteFile(t, deletePath, 1, 128) + task := FileScanTask{ + File: newLazyDataFile(t, dataPath), + DeleteFiles: []iceberg.DataFile{deleteFile}, + } + loader := newLazyPositionDeleteLoader(fs, []FileScanTask{task}) + + const callers = 8 + results := make(chan positionDeletes, callers) + errs := make(chan error, callers) + var wg sync.WaitGroup + wg.Add(callers) + for range callers { + go func() { + defer wg.Done() + deletes, err := loader.load(ctx, task) + results <- deletes + errs <- err + }() + } + wg.Wait() + close(results) + close(errs) + + for err := range errs { + require.NoError(t, err) + } + for deletes := range results { + require.Len(t, deletes, 1) + assert.Equal(t, []int64{7}, int64Values(deletes[0])) + } + assert.Equal(t, int64(1), fs.opens.Load(), "concurrent users must share one delete-file read") + + loader.release() +} + +func TestLazyPositionDeleteLoaderCachesErrors(t *testing.T) { + fs := &countingOpenMemFS{MemFS: iceio.NewMemFS()} + deleteFile := newPosDeleteFile(t, "mem://bucket/deletes/missing.parquet", 1, 128) + task := FileScanTask{ + File: newLazyDataFile(t, "mem://bucket/data/a.parquet"), + DeleteFiles: []iceberg.DataFile{deleteFile}, + } + loader := newLazyPositionDeleteLoader(fs, []FileScanTask{task}) + + first, err := loader.load(context.Background(), task) + require.Error(t, err) + assert.Nil(t, first) + assert.Equal(t, int64(1), fs.opens.Load()) + + second, secondErr := loader.load(context.Background(), task) + require.Error(t, secondErr) + assert.Nil(t, second) + assert.ErrorIs(t, secondErr, err) + assert.Equal(t, int64(1), fs.opens.Load(), "a failed delete file must not be retried by other tasks") + + loader.release() +} + +func TestLazyPositionDeleteLoaderCachesCancellation(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + ctx, cancel := context.WithCancel(compute.WithAllocator(t.Context(), mem)) + cancel() + + fs := &countingOpenMemFS{MemFS: iceio.NewMemFS()} + deletePath := "mem://bucket/deletes/cancelled.parquet" + writePosDeleteParquetToMemFS(t, fs.MemFS, deletePath, `[{"file_path":"mem://bucket/data/a.parquet","pos":1}]`) + deleteFile := newPosDeleteFile(t, deletePath, 1, 128) + task := FileScanTask{ + File: newLazyDataFile(t, "mem://bucket/data/a.parquet"), + DeleteFiles: []iceberg.DataFile{deleteFile}, + } + loader := newLazyPositionDeleteLoader(fs, []FileScanTask{task}) + + _, firstErr := loader.load(ctx, task) + require.ErrorIs(t, firstErr, context.Canceled) + _, secondErr := loader.load(context.Background(), task) + require.ErrorIs(t, secondErr, context.Canceled) + assert.Equal(t, int64(1), fs.opens.Load(), "cancellation must not cause a second read") + + loader.release() +} + +func TestArrowScanDefersPositionDeleteReadsUntilIteration(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + ctx := compute.WithAllocator(t.Context(), mem) + + schema := iceberg.NewSchema(1, iceberg.NestedField{ + ID: 1, Name: "value", Type: iceberg.PrimitiveTypes.Int64, + }) + metadata, err := NewMetadata(schema, iceberg.UnpartitionedSpec, + UnsortedSortOrder, "mem://bucket/table", nil) + require.NoError(t, err) + + const deletePath = "mem://bucket/deletes/one.parquet" + memFS := &countingOpenMemFS{MemFS: iceio.NewMemFS()} + writePosDeleteParquetToMemFS(t, memFS.MemFS, deletePath, + `[{"file_path":"mem://bucket/data/missing.parquet","pos":1}]`) + deleteFile := newPosDeleteFile(t, deletePath, 1, 128) + task := FileScanTask{ + File: newLazyDataFile(t, "mem://bucket/data/missing.parquet"), + DeleteFiles: []iceberg.DataFile{deleteFile}, + } + scan := &arrowScan{ + metadata: metadata, + fs: memFS, + scanSchema: schema, + projectedSchema: schema, + boundRowFilter: iceberg.AlwaysTrue{}, + rowLimit: -1, + concurrency: 1, + } + + _, records, err := scan.GetRecords(ctx, []FileScanTask{task}) + require.NoError(t, err) + assert.Zero(t, memFS.opens.Load(), "GetRecords must not read position deletes") + + var iterErr error + for record, err := range records { + if record != nil { + record.Release() + } + iterErr = err + break + } + require.Error(t, iterErr, "iteration should reach the missing data file after loading its delete") + assert.Equal(t, int64(2), memFS.opens.Load(), "the first task should open one delete and one data file") +} + +func TestLazyPositionDeleteLoaderReleasesChunksWhenIteratorStops(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + + loader := &lazyPositionDeleteLoader{files: map[string]*lazyPositionDeleteFile{ + "mem://bucket/deletes/one.parquet": { + deletes: map[string]*arrow.Chunked{ + "mem://bucket/data/a.parquet": chunkedPosDelete(t, mem, []int64{1}), + }, + }, + }} + + batch := checkedInt64RecordBatch(mem, 1) + records := make(chan enumeratedRecord, 1) + records <- enumeratedRecord{ + Record: internal.Enumerated[arrow.RecordBatch]{ + Value: batch, + Index: 0, + Last: true, + }, + Task: internal.Enumerated[FileScanTask]{Index: 0, Last: true}, + } + close(records) + + ctx, cancel := context.WithCancelCause(context.Background()) + itr := createIteratorWithCleanup(ctx, 1, records, nil, cancel, 0, loader.release) + for record, err := range itr { + require.NoError(t, err) + record.Release() + break + } +} From 79c740e6b8484056d7ad01fa4cf9bd4c135f9b90 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Fri, 28 Aug 2026 13:52:05 +0200 Subject: [PATCH 02/10] chore(table): satisfy lint Signed-off-by: Minh Vu --- table/arrow_scanner_lazy_delete_regression_test.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/table/arrow_scanner_lazy_delete_regression_test.go b/table/arrow_scanner_lazy_delete_regression_test.go index 04b69ed2e..c6a1da164 100644 --- a/table/arrow_scanner_lazy_delete_regression_test.go +++ b/table/arrow_scanner_lazy_delete_regression_test.go @@ -238,6 +238,7 @@ func TestArrowScanDefersPositionDeleteReadsUntilIteration(t *testing.T) { record.Release() } iterErr = err + break } require.Error(t, iterErr, "iteration should reach the missing data file after loading its delete") @@ -273,6 +274,7 @@ func TestLazyPositionDeleteLoaderReleasesChunksWhenIteratorStops(t *testing.T) { for record, err := range itr { require.NoError(t, err) record.Release() + break } } From 5437fbc98374103888ed26203cf4f5063f7659a6 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Fri, 28 Aug 2026 15:43:10 +0200 Subject: [PATCH 03/10] perf(table): avoid single-delete dedup allocation Signed-off-by: Minh Vu --- table/arrow_scanner.go | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/table/arrow_scanner.go b/table/arrow_scanner.go index 16d8ba505..43651a05c 100644 --- a/table/arrow_scanner.go +++ b/table/arrow_scanner.go @@ -188,17 +188,24 @@ func (l *lazyPositionDeleteLoader) load(ctx context.Context, task FileScanTask) targetPath := task.File.FilePath() deletes := make(positionDeletes, 0, len(task.DeleteFiles)) - seen := make(map[string]struct{}, len(task.DeleteFiles)) + // Most scan tasks carry one positional delete file. Avoid allocating a + // deduplication map unless there can actually be duplicate entries. + var seen map[string]struct{} + if len(task.DeleteFiles) > 1 { + seen = make(map[string]struct{}, len(task.DeleteFiles)) + } for _, deleteFile := range task.DeleteFiles { if deleteFile.ContentType() != iceberg.EntryContentPosDeletes { continue } path := deleteFile.FilePath() - if _, ok := seen[path]; ok { - continue + if seen != nil { + if _, ok := seen[path]; ok { + continue + } + seen[path] = struct{}{} } - seen[path] = struct{}{} cached, ok := l.files[path] if !ok { From dd27e8c9978da082491678608c0a7eddd6c1fd89 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Fri, 28 Aug 2026 22:36:28 +0200 Subject: [PATCH 04/10] fix(table): tear down lazy scan workers on cancellation Signed-off-by: Minh Vu --- table/arrow_scanner.go | 18 ++-- ...row_scanner_lazy_delete_regression_test.go | 99 +++++++++++++++++++ table/internal/utils.go | 16 +++ 3 files changed, 127 insertions(+), 6 deletions(-) diff --git a/table/arrow_scanner.go b/table/arrow_scanner.go index 43651a05c..1c7a08693 100644 --- a/table/arrow_scanner.go +++ b/table/arrow_scanner.go @@ -1824,7 +1824,7 @@ func createIteratorWithCleanup(ctx context.Context, numWorkers uint, records <-c return batch.Task.Index < 0 } - sequenced := tblutils.MakeSequencedChan(uint(numWorkers), records, + sequenced := tblutils.MakeSequencedChanWithDiscard(uint(numWorkers), records, func(left, right *enumeratedRecord) bool { switch { case isBeforeAny(*left): @@ -1850,7 +1850,11 @@ func createIteratorWithCleanup(ctx context.Context, numWorkers uint, records <-c return next.Task.Index == prev.Task.Index+1 && prev.Record.Last && next.Record.Index == 0 } - }, enumeratedRecord{Task: tblutils.Enumerated[FileScanTask]{Index: -1}}) + }, enumeratedRecord{Task: tblutils.Enumerated[FileScanTask]{Index: -1}}, func(rec enumeratedRecord) { + if rec.Record.Value != nil { + rec.Record.Value.Release() + } + }) totalRowCount := int64(0) @@ -1973,6 +1977,12 @@ func (as *arrowScan) recordBatchesFromTasksAndDeletes(ctx context.Context, tasks } go func() { + defer func() { + close(taskChan) + wg.Wait() + close(records) + }() + for i, t := range tasks { select { case <-scanCtx.Done(): @@ -1982,10 +1992,6 @@ func (as *arrowScan) recordBatchesFromTasksAndDeletes(ctx context.Context, tasks }: } } - close(taskChan) - - wg.Wait() - close(records) }() var cleanup func() diff --git a/table/arrow_scanner_lazy_delete_regression_test.go b/table/arrow_scanner_lazy_delete_regression_test.go index c6a1da164..fe5433948 100644 --- a/table/arrow_scanner_lazy_delete_regression_test.go +++ b/table/arrow_scanner_lazy_delete_regression_test.go @@ -19,9 +19,11 @@ package table import ( "context" + "errors" "sync" "sync/atomic" "testing" + "time" "github.com/apache/arrow-go/v18/arrow" "github.com/apache/arrow-go/v18/arrow/compute" @@ -278,3 +280,100 @@ func TestLazyPositionDeleteLoaderReleasesChunksWhenIteratorStops(t *testing.T) { break } } + +func TestArrowScanPreCancelledIteratorTearsDownProducer(t *testing.T) { + ctx, cancel := context.WithCancelCause(context.Background()) + cancel(context.Canceled) + + scan := &arrowScan{concurrency: 1, rowLimit: -1} + tasks := []FileScanTask{{ + File: newLazyDataFile(t, "mem://bucket/data/pre-cancelled.parquet"), + }} + records := scan.recordBatchesFromTasksAndDeletes(ctx, tasks, nil, nil, nil, nil) + + done := make(chan error, 1) + go func() { + var iterErr error + for _, err := range records { + if err != nil { + iterErr = err + } + } + done <- iterErr + }() + + select { + case err := <-done: + require.ErrorIs(t, err, context.Canceled) + case <-time.After(500 * time.Millisecond): + t.Fatal("pre-cancelled scan iterator did not terminate") + } +} + +func TestCreateIteratorReleasesOutOfOrderBatchAfterError(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + + expectedErr := errors.New("position delete load failed") + records := make(chan enumeratedRecord, 2) + records <- enumeratedRecord{ + Record: internal.Enumerated[arrow.RecordBatch]{ + Value: checkedInt64RecordBatch(mem, 1), + Index: 0, + Last: true, + }, + Task: internal.Enumerated[FileScanTask]{Index: 1, Last: true}, + } + records <- enumeratedRecord{ + Task: internal.Enumerated[FileScanTask]{Index: 0}, + Err: expectedErr, + } + close(records) + + ctx, cancel := context.WithCancelCause(context.Background()) + var gotErr error + for _, err := range createIteratorWithCleanup(ctx, 2, records, nil, cancel, 0, nil) { + gotErr = err + } + + require.ErrorIs(t, gotErr, expectedErr) +} + +func TestCreateIteratorStopsWhenConsumerStops(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + + records := make(chan enumeratedRecord, 1) + records <- enumeratedRecord{ + Record: internal.Enumerated[arrow.RecordBatch]{ + Value: checkedInt64RecordBatch(mem, 1), + Index: 0, + Last: true, + }, + Task: internal.Enumerated[FileScanTask]{Index: 0, Last: true}, + } + + ctx, cancel := context.WithCancelCause(context.Background()) + defer func() { cancel(nil) }() + closed := make(chan struct{}) + go func() { + <-ctx.Done() + close(records) + }() + + go func() { + for record, err := range createIteratorWithCleanup(ctx, 1, records, nil, cancel, 0, nil) { + if err == nil && record != nil { + record.Release() + } + break + } + close(closed) + }() + + select { + case <-closed: + case <-time.After(500 * time.Millisecond): + t.Fatal("iterator did not stop after consumer termination") + } +} diff --git a/table/internal/utils.go b/table/internal/utils.go index 5fb7e6164..743985db8 100644 --- a/table/internal/utils.go +++ b/table/internal/utils.go @@ -83,11 +83,27 @@ func (pq *pqueue[T]) Pop() any { // based on the comesAfter and isNext functions. The values are read in from // the provided source and then re-ordered before being sent to the output. func MakeSequencedChan[T any](bufferSize uint, source <-chan T, comesAfter, isNext func(a, b *T) bool, initial T) <-chan T { + return MakeSequencedChanWithDiscard(bufferSize, source, comesAfter, isNext, initial, nil) +} + +// MakeSequencedChanWithDiscard is MakeSequencedChan with a callback for values +// that remain in the reorder queue when the source is closed. Values already +// sent to the output are not passed to discard. +func MakeSequencedChanWithDiscard[T any](bufferSize uint, source <-chan T, comesAfter, isNext func(a, b *T) bool, initial T, discard func(T)) <-chan T { pq := pqueue[T]{queue: make([]*T, 0), compare: comesAfter} heap.Init(&pq) previous, out := &initial, make(chan T, bufferSize) go func() { defer close(out) + defer func() { + if discard == nil { + return + } + + for _, val := range pq.queue { + discard(*val) + } + }() for val := range source { heap.Push(&pq, &val) for pq.Len() > 0 && isNext(previous, pq.queue[0]) { From fb3a7d3c6f07af30d1b0ffaa7fccbe7b6fb419e1 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Fri, 28 Aug 2026 22:54:04 +0200 Subject: [PATCH 05/10] chore(table): satisfy nlreturn Signed-off-by: Minh Vu --- table/arrow_scanner_lazy_delete_regression_test.go | 1 + 1 file changed, 1 insertion(+) diff --git a/table/arrow_scanner_lazy_delete_regression_test.go b/table/arrow_scanner_lazy_delete_regression_test.go index fe5433948..cd36cf67e 100644 --- a/table/arrow_scanner_lazy_delete_regression_test.go +++ b/table/arrow_scanner_lazy_delete_regression_test.go @@ -366,6 +366,7 @@ func TestCreateIteratorStopsWhenConsumerStops(t *testing.T) { if err == nil && record != nil { record.Release() } + break } close(closed) From 02346d6f2e0496b248d63fae410d97d67f4bd030 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Sun, 30 Aug 2026 23:28:52 +0200 Subject: [PATCH 06/10] fix(schema): avoid copying lazy caches during JSON encoding Signed-off-by: Minh Vu --- schema.go | 3 +-- schema_test.go | 46 ++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 47 insertions(+), 2 deletions(-) diff --git a/schema.go b/schema.go index 24b0ff4b1..4330103dc 100644 --- a/schema.go +++ b/schema.go @@ -346,8 +346,7 @@ func (s *Schema) MarshalJSON() ([]byte, error) { type Alias Schema - aliasCopy := *(*Alias)(s) - aliasCopy.IdentifierFieldIDs = ids + aliasCopy := Alias{ID: s.ID, IdentifierFieldIDs: ids} return json.Marshal(struct { Type string `json:"type"` diff --git a/schema_test.go b/schema_test.go index 7fa16e345..6d794ab71 100644 --- a/schema_test.go +++ b/schema_test.go @@ -24,6 +24,7 @@ import ( "path/filepath" "runtime" "strings" + "sync" "testing" "github.com/apache/iceberg-go" @@ -2293,3 +2294,48 @@ func TestVisitGeoSchemaWithSchemaVisitorPerPrimitiveType(t *testing.T) { assert.Equal(t, 1, v.geometryCalls) assert.Equal(t, 1, v.geographyCalls) } + +func TestSchemaMarshalJSONConcurrentLazyLookups(t *testing.T) { + for range 32 { + schema := iceberg.NewSchemaWithIdentifiers(17, nil, + iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64, Required: true}, + iceberg.NestedField{ID: 2, Name: "data", Type: iceberg.PrimitiveTypes.String}, + ) + start := make(chan struct{}) + var wg sync.WaitGroup + for range 8 { + wg.Go(func() { + <-start + for range 8 { + _, err := json.Marshal(schema) + assert.NoError(t, err) + } + }) + wg.Go(func() { + <-start + _, found := schema.FindFieldByID(1) + assert.True(t, found) + _, found = schema.FindFieldByName("data") + assert.True(t, found) + _, found = schema.FindFieldByNameCaseInsensitive("DATA") + assert.True(t, found) + name, found := schema.FindColumnName(2) + assert.True(t, found) + assert.Equal(t, "data", name) + }) + } + close(start) + wg.Wait() + + data, err := json.Marshal(schema) + require.NoError(t, err) + assert.JSONEq(t, `{ + "type": "struct", "schema-id": 17, "identifier-field-ids": [], + "fields": [ + {"id": 1, "name": "id", "type": "long", "required": true}, + {"id": 2, "name": "data", "type": "string", "required": false} + ] + }`, string(data)) + assert.Nil(t, schema.IdentifierFieldIDs) + } +} From 62ae710507d791d9cad639613803736d29fa6a07 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Mon, 31 Aug 2026 10:29:35 +0200 Subject: [PATCH 07/10] fix(table): finish lazy position delete review fixes --- table/arrow_scanner.go | 16 ++- ...row_scanner_lazy_delete_regression_test.go | 115 +++++++++++++++++- table/scanner.go | 3 + 3 files changed, 132 insertions(+), 2 deletions(-) diff --git a/table/arrow_scanner.go b/table/arrow_scanner.go index 1c7a08693..5ba8f92c5 100644 --- a/table/arrow_scanner.go +++ b/table/arrow_scanner.go @@ -73,6 +73,8 @@ func releasePerFilePosDeletes(deletesPerFile perFilePosDeletes) { } } +// readAllDeleteFiles is retained for the eager-path benchmark and regression +// tests; scans use lazyPositionDeleteLoader instead. func readAllDeleteFiles(ctx context.Context, fs iceio.IO, tasks []FileScanTask, concurrency int) (perFilePosDeletes, error) { deletesPerFile := make(perFilePosDeletes) uniqueDeletes := make(map[string]iceberg.DataFile) @@ -144,6 +146,9 @@ func readAllDeleteFiles(ctx context.Context, fs iceio.IO, tasks []FileScanTask, // The grouped Arrow chunks are owned by the loader until release. The iterator // calls release after all workers have stopped, which keeps shared chunks alive // while multiple tasks use them and also covers early iterator termination. +// The loader has the lifetime of exactly one scan. Each file's first load also +// locks in its result, including context errors, for every caller; a loader +// must not be reused for a retry with a different context. type lazyPositionDeleteLoader struct { fs iceio.IO files map[string]*lazyPositionDeleteFile @@ -222,6 +227,8 @@ func (l *lazyPositionDeleteLoader) load(ctx context.Context, task FileScanTask) // in case a future reader returns partial Arrow ownership. releasePosDeletes(cached.deletes) cached.deletes = nil + cached.err = fmt.Errorf("read position deletes from %s: %w", + cached.dataFile.FilePath(), cached.err) } }) if cached.err != nil { @@ -1955,7 +1962,10 @@ func (as *arrowScan) recordBatchesFromTasksAndDeletes(ctx context.Context, tasks var err error positionalDeletes, err = positionDeleteLoader.load(scanCtx, task.Value) if err != nil { - records <- enumeratedRecord{Task: task, Err: err} + select { + case records <- enumeratedRecord{Task: task, Err: err}: + case <-scanCtx.Done(): + } cancel(err) return @@ -2003,6 +2013,10 @@ func (as *arrowScan) recordBatchesFromTasksAndDeletes(ctx context.Context, tasks } } +// GetRecords prepares the projected Arrow schema and a single-use record +// iterator. Positional-delete files are opened and read during iteration, so +// errors from those files are returned by the iterator; deletion-vector and +// equality-delete errors are returned before the iterator is created. func (as *arrowScan) GetRecords(ctx context.Context, tasks []FileScanTask) (*arrow.Schema, iter.Seq2[arrow.RecordBatch, error], error) { var err error as.useLargeTypes, err = strconv.ParseBool(as.options.Get(ScanOptionArrowUseLargeTypes, "false")) diff --git a/table/arrow_scanner_lazy_delete_regression_test.go b/table/arrow_scanner_lazy_delete_regression_test.go index cd36cf67e..86dded818 100644 --- a/table/arrow_scanner_lazy_delete_regression_test.go +++ b/table/arrow_scanner_lazy_delete_regression_test.go @@ -20,14 +20,19 @@ package table import ( "context" "errors" + "fmt" + "strings" "sync" "sync/atomic" "testing" "time" "github.com/apache/arrow-go/v18/arrow" + "github.com/apache/arrow-go/v18/arrow/array" "github.com/apache/arrow-go/v18/arrow/compute" "github.com/apache/arrow-go/v18/arrow/memory" + "github.com/apache/arrow-go/v18/parquet" + "github.com/apache/arrow-go/v18/parquet/pqarrow" "github.com/apache/iceberg-go" iceio "github.com/apache/iceberg-go/io" "github.com/apache/iceberg-go/table/internal" @@ -37,11 +42,16 @@ import ( type countingOpenMemFS struct { *iceio.MemFS - opens atomic.Int64 + opens atomic.Int64 + trackedPath string + trackedOpens atomic.Int64 } func (f *countingOpenMemFS) Open(name string) (iceio.File, error) { f.opens.Add(1) + if name == f.trackedPath { + f.trackedOpens.Add(1) + } return f.MemFS.Open(name) } @@ -57,6 +67,42 @@ func newLazyDataFile(t *testing.T, path string) iceberg.DataFile { return builder.Build() } +func writeLazyDataParquetToMemFS(t *testing.T, fs *iceio.MemFS, path string, start, count int) iceberg.DataFile { + t.Helper() + + dataSchema := arrow.NewSchema([]arrow.Field{{ + Name: "value", + Type: arrow.PrimitiveTypes.Int64, + Nullable: false, + Metadata: arrow.MetadataFrom(map[string]string{ArrowParquetFieldIDKey: "1"}), + }}, nil) + bldr := array.NewInt64Builder(memory.DefaultAllocator) + defer bldr.Release() + for i := range count { + bldr.Append(int64(start + i)) + } + values := bldr.NewArray() + defer values.Release() + record := array.NewRecordBatch(dataSchema, []arrow.Array{values}, int64(count)) + defer record.Release() + tbl := array.NewTableFromRecords(dataSchema, []arrow.RecordBatch{record}) + defer tbl.Release() + + file, err := fs.Create(path) + require.NoError(t, err) + require.NoError(t, pqarrow.WriteTable(tbl, file, int64(count), + parquet.NewWriterProperties(parquet.WithStats(true)), + pqarrow.DefaultWriterProps())) + require.NoError(t, file.Close()) + + builder, err := iceberg.NewDataFileBuilder( + *iceberg.UnpartitionedSpec, iceberg.EntryContentData, + path, iceberg.ParquetFile, nil, nil, nil, int64(count), 128) + require.NoError(t, err) + + return builder.Build() +} + func TestLazyPositionDeleteLoaderDefersReadsAndSharesResults(t *testing.T) { mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) defer mem.AssertSize(t, 0) @@ -163,6 +209,7 @@ func TestLazyPositionDeleteLoaderCachesErrors(t *testing.T) { first, err := loader.load(context.Background(), task) require.Error(t, err) assert.Nil(t, first) + assert.Contains(t, err.Error(), deleteFile.FilePath()) assert.Equal(t, int64(1), fs.opens.Load()) second, secondErr := loader.load(context.Background(), task) @@ -247,6 +294,72 @@ func TestArrowScanDefersPositionDeleteReadsUntilIteration(t *testing.T) { assert.Equal(t, int64(2), memFS.opens.Load(), "the first task should open one delete and one data file") } +func TestArrowScanReleasesLazyPositionDeletesOnEarlyStop(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + ctx := compute.WithAllocator(t.Context(), mem) + + schema := iceberg.NewSchema(1, iceberg.NestedField{ + ID: 1, Name: "value", Type: iceberg.PrimitiveTypes.Int64, + }) + metadata, err := NewMetadata(schema, iceberg.UnpartitionedSpec, + UnsortedSortOrder, "mem://bucket/table", nil) + require.NoError(t, err) + + const ( + deletePath = "mem://bucket/deletes/shared.parquet" + taskCount = 8 + rowCount = 4 + ) + memFS := &countingOpenMemFS{ + MemFS: iceio.NewMemFS(), + trackedPath: deletePath, + } + tasks := make([]FileScanTask, taskCount) + var deleteJSON strings.Builder + deleteJSON.WriteByte('[') + for i := range tasks { + dataPath := fmt.Sprintf("mem://bucket/data/data-%d.parquet", i) + tasks[i].File = writeLazyDataParquetToMemFS(t, memFS.MemFS, + dataPath, i*rowCount, rowCount) + if i > 0 { + deleteJSON.WriteByte(',') + } + fmt.Fprintf(&deleteJSON, `{"file_path":"%s","pos":0}`, dataPath) + } + deleteJSON.WriteByte(']') + writePosDeleteParquetToMemFS(t, memFS.MemFS, deletePath, deleteJSON.String()) + deleteFile := newPosDeleteFile(t, deletePath, taskCount, 128) + for i := range tasks { + tasks[i].DeleteFiles = []iceberg.DataFile{deleteFile} + } + + scan := &arrowScan{ + metadata: metadata, + fs: memFS, + scanSchema: schema, + projectedSchema: schema, + boundRowFilter: iceberg.AlwaysTrue{}, + rowLimit: -1, + concurrency: 4, + } + _, records, err := scan.GetRecords(ctx, tasks) + require.NoError(t, err) + + var batches int + for record, err := range records { + require.NoError(t, err) + require.NotNil(t, record) + record.Release() + batches++ + break + } + + assert.Equal(t, 1, batches, "the test must stop after the first batch") + assert.Equal(t, int64(1), memFS.trackedOpens.Load(), + "all workers must share one lazily-loaded positional-delete file") +} + func TestLazyPositionDeleteLoaderReleasesChunksWhenIteratorStops(t *testing.T) { mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) defer mem.AssertSize(t, 0) diff --git a/table/scanner.go b/table/scanner.go index f5599d932..982d4702a 100644 --- a/table/scanner.go +++ b/table/scanner.go @@ -1467,6 +1467,9 @@ func (scan *Scan) ToArrowRecords(ctx context.Context) (*arrow.Schema, iter.Seq2[ // ReadTasks reads Arrow records from a specific set of FileScanTasks, applying the // scan's projection, per-task residual filters, and positional delete handling. This // is useful when the caller has already planned or selected specific tasks to read. +// Positional-delete read errors are returned by the iterator during iteration; +// deletion-vector and equality-delete read errors are returned by ReadTasks +// before it returns an iterator. The returned iterator is single-use. func (scan *Scan) ReadTasks(ctx context.Context, tasks []FileScanTask) (*arrow.Schema, iter.Seq2[arrow.RecordBatch, error], error) { if atomic.LoadUint32(&scan.closed) != 0 { return nil, nil, fmt.Errorf("%w: scan is closed", ErrInvalidOperation) From 4a3a548064ded96da090d410c0850498647ebb48 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Mon, 31 Aug 2026 10:43:00 +0200 Subject: [PATCH 08/10] chore(table): satisfy nlreturn --- table/arrow_scanner_lazy_delete_regression_test.go | 1 + 1 file changed, 1 insertion(+) diff --git a/table/arrow_scanner_lazy_delete_regression_test.go b/table/arrow_scanner_lazy_delete_regression_test.go index 86dded818..f7ca2aa03 100644 --- a/table/arrow_scanner_lazy_delete_regression_test.go +++ b/table/arrow_scanner_lazy_delete_regression_test.go @@ -352,6 +352,7 @@ func TestArrowScanReleasesLazyPositionDeletesOnEarlyStop(t *testing.T) { require.NotNil(t, record) record.Release() batches++ + break } From 66cedb31e64d21cfde67d9a7d72755e399437c2e Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Fri, 4 Sep 2026 14:32:00 +0200 Subject: [PATCH 09/10] fix(table): address lazy position delete review feedback --- table/arrow_scanner.go | 21 ++++- table/arrow_scanner_lazy_delete_bench_test.go | 6 -- ...row_scanner_lazy_delete_regression_test.go | 93 ++++++++++++++++++- table/scanner.go | 8 +- 4 files changed, 114 insertions(+), 14 deletions(-) diff --git a/table/arrow_scanner.go b/table/arrow_scanner.go index 7add57f73..ae8780891 100644 --- a/table/arrow_scanner.go +++ b/table/arrow_scanner.go @@ -19,6 +19,7 @@ package table import ( "context" + "errors" "fmt" "io" "iter" @@ -26,6 +27,7 @@ import ( "strconv" "strings" "sync" + "sync/atomic" "github.com/apache/arrow-go/v18/arrow" "github.com/apache/arrow-go/v18/arrow/array" @@ -76,8 +78,9 @@ func releasePerFilePosDeletes(deletesPerFile perFilePosDeletes) { } } -// readAllDeleteFiles is retained for the eager-path benchmark and regression -// tests; scans use lazyPositionDeleteLoader instead. +// readAllDeleteFiles reads every referenced positional-delete file up front. It +// remains the path used by Transaction.makePositionDeleteRecordsForFilter and by +// the eager-vs-lazy benchmark; arrowScan.GetRecords uses lazyPositionDeleteLoader. func readAllDeleteFiles(ctx context.Context, fs iceio.IO, tasks []FileScanTask, concurrency int) (perFilePosDeletes, error) { deletesPerFile := make(perFilePosDeletes) uniqueDeletes := make(map[string]iceberg.DataFile) @@ -157,8 +160,11 @@ type lazyPositionDeleteLoader struct { files map[string]*lazyPositionDeleteFile releaseOnce sync.Once + released atomic.Bool } +var errPositionDeleteLoaderReleased = errors.New("position delete loader already released") + type lazyPositionDeleteFile struct { dataFile iceberg.DataFile @@ -190,6 +196,10 @@ func newLazyPositionDeleteLoader(fs iceio.IO, tasks []FileScanTask) *lazyPositio } func (l *lazyPositionDeleteLoader) load(ctx context.Context, task FileScanTask) (positionDeletes, error) { + if l.released.Load() { + return nil, errPositionDeleteLoaderReleased + } + if len(task.DeleteFiles) == 0 { return nil, nil } @@ -252,6 +262,7 @@ func (l *lazyPositionDeleteLoader) release() { } l.releaseOnce.Do(func() { + l.released.Store(true) for _, cached := range l.files { releasePosDeletes(cached.deletes) cached.deletes = nil @@ -2280,8 +2291,10 @@ func (as *arrowScan) recordBatchesFromTasksAndDeletes(ctx context.Context, tasks // GetRecords prepares the projected Arrow schema and a single-use record // iterator. Positional- and equality-delete files are opened and read during -// iteration, so errors from those files are returned by the iterator; -// deletion-vector errors are returned before the iterator is created. +// iteration. Errors from those files are delivered through the iterator only if +// iteration reaches the task that encounters the error; a row limit or early +// termination may finish the scan before the error is observed. Deletion-vector +// errors are returned before the iterator is created. func (as *arrowScan) GetRecords(ctx context.Context, tasks []FileScanTask) (*arrow.Schema, iter.Seq2[arrow.RecordBatch, error], error) { var err error as.useLargeTypes, err = strconv.ParseBool(as.options.Get(ScanOptionArrowUseLargeTypes, "false")) diff --git a/table/arrow_scanner_lazy_delete_bench_test.go b/table/arrow_scanner_lazy_delete_bench_test.go index da3a22f92..b25243a87 100644 --- a/table/arrow_scanner_lazy_delete_bench_test.go +++ b/table/arrow_scanner_lazy_delete_bench_test.go @@ -126,8 +126,6 @@ func BenchmarkLazyPositionDeleteLoading(b *testing.B) { } releasePerFilePosDeletes(deletes) } - b.ReportMetric(lazyPositionDeleteBenchmarkFileCount, "delete_files/op") - b.ReportMetric(lazyPositionDeleteBenchmarkFileCount, "delete_files_before_first/op") }) b.Run("lazy_unread_iterator", func(b *testing.B) { @@ -141,8 +139,6 @@ func BenchmarkLazyPositionDeleteLoading(b *testing.B) { } loader.release() } - b.ReportMetric(0, "delete_files/op") - b.ReportMetric(0, "delete_files_before_first/op") }) b.Run("lazy_first_task", func(b *testing.B) { @@ -159,7 +155,5 @@ func BenchmarkLazyPositionDeleteLoading(b *testing.B) { } loader.release() } - b.ReportMetric(1, "delete_files/op") - b.ReportMetric(1, "delete_files_before_first/op") }) } diff --git a/table/arrow_scanner_lazy_delete_regression_test.go b/table/arrow_scanner_lazy_delete_regression_test.go index f7ca2aa03..ba0dc0bf8 100644 --- a/table/arrow_scanner_lazy_delete_regression_test.go +++ b/table/arrow_scanner_lazy_delete_regression_test.go @@ -45,6 +45,8 @@ type countingOpenMemFS struct { opens atomic.Int64 trackedPath string trackedOpens atomic.Int64 + blockPath string + unblock <-chan struct{} } func (f *countingOpenMemFS) Open(name string) (iceio.File, error) { @@ -52,6 +54,9 @@ func (f *countingOpenMemFS) Open(name string) (iceio.File, error) { if name == f.trackedPath { f.trackedOpens.Add(1) } + if name == f.blockPath && f.unblock != nil { + <-f.unblock + } return f.MemFS.Open(name) } @@ -209,7 +214,7 @@ func TestLazyPositionDeleteLoaderCachesErrors(t *testing.T) { first, err := loader.load(context.Background(), task) require.Error(t, err) assert.Nil(t, first) - assert.Contains(t, err.Error(), deleteFile.FilePath()) + assert.ErrorContains(t, err, "read position deletes from "+deleteFile.FilePath()) assert.Equal(t, int64(1), fs.opens.Load()) second, secondErr := loader.load(context.Background(), task) @@ -221,6 +226,22 @@ func TestLazyPositionDeleteLoaderCachesErrors(t *testing.T) { loader.release() } +func TestLazyPositionDeleteLoaderWrapsErrorsWithoutFilePath(t *testing.T) { + fs := &countingOpenMemFS{MemFS: iceio.NewMemFS()} + deletePath := "mem://bucket/deletes/corrupt.parquet" + require.NoError(t, fs.WriteFile(deletePath, []byte("not a parquet file"))) + task := FileScanTask{ + File: newLazyDataFile(t, "mem://bucket/data/a.parquet"), + DeleteFiles: []iceberg.DataFile{newPosDeleteFile(t, deletePath, 1, 128)}, + } + loader := newLazyPositionDeleteLoader(fs, []FileScanTask{task}) + + _, err := loader.load(context.Background(), task) + require.Error(t, err) + assert.ErrorContains(t, err, "read position deletes from "+deletePath) + loader.release() +} + func TestLazyPositionDeleteLoaderCachesCancellation(t *testing.T) { mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) defer mem.AssertSize(t, 0) @@ -361,6 +382,64 @@ func TestArrowScanReleasesLazyPositionDeletesOnEarlyStop(t *testing.T) { "all workers must share one lazily-loaded positional-delete file") } +func TestArrowScanRowLimitStopsBeforeLazyPositionDeleteError(t *testing.T) { + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + ctx := compute.WithAllocator(t.Context(), mem) + + schema := iceberg.NewSchema(1, iceberg.NestedField{ + ID: 1, Name: "value", Type: iceberg.PrimitiveTypes.Int64, + }) + metadata, err := NewMetadata(schema, iceberg.UnpartitionedSpec, + UnsortedSortOrder, "mem://bucket/table", nil) + require.NoError(t, err) + + const corruptDeletePath = "mem://bucket/deletes/corrupt.parquet" + unblock := make(chan struct{}) + memFS := &countingOpenMemFS{ + MemFS: iceio.NewMemFS(), + blockPath: corruptDeletePath, + unblock: unblock, + } + const taskCount = 4 + tasks := make([]FileScanTask, taskCount) + for i := range tasks { + dataPath := fmt.Sprintf("mem://bucket/data/data-%d.parquet", i) + tasks[i].File = writeLazyDataParquetToMemFS(t, memFS.MemFS, dataPath, i, 4) + } + require.NoError(t, memFS.WriteFile(corruptDeletePath, []byte("not a parquet file"))) + tasks[taskCount-1].DeleteFiles = []iceberg.DataFile{ + newPosDeleteFile(t, corruptDeletePath, 1, 128), + } + + scan := &arrowScan{ + metadata: metadata, + fs: memFS, + scanSchema: schema, + projectedSchema: schema, + boundRowFilter: iceberg.AlwaysTrue{}, + rowLimit: 1, + concurrency: 2, + } + _, records, err := scan.GetRecords(ctx, tasks) + require.NoError(t, err) + + var releaseBlockedDelete sync.Once + release := func() { releaseBlockedDelete.Do(func() { close(unblock) }) } + defer release() + var rows int64 + for record, err := range records { + require.NoError(t, err) + require.NotNil(t, record) + rows += record.NumRows() + record.Release() + release() + } + + assert.Equal(t, int64(1), rows, + "a row limit may finish before a later task's positional-delete error") +} + func TestLazyPositionDeleteLoaderReleasesChunksWhenIteratorStops(t *testing.T) { mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) defer mem.AssertSize(t, 0) @@ -393,6 +472,18 @@ func TestLazyPositionDeleteLoaderReleasesChunksWhenIteratorStops(t *testing.T) { break } + + for _, cached := range loader.files { + assert.Nil(t, cached.deletes, "release must clear cached delete chunks") + } +} + +func TestLazyPositionDeleteLoaderRejectsLoadAfterRelease(t *testing.T) { + loader := &lazyPositionDeleteLoader{} + loader.release() + + _, err := loader.load(context.Background(), FileScanTask{}) + require.ErrorIs(t, err, errPositionDeleteLoaderReleased) } func TestArrowScanPreCancelledIteratorTearsDownProducer(t *testing.T) { diff --git a/table/scanner.go b/table/scanner.go index 2ce7e223d..dd8c74942 100644 --- a/table/scanner.go +++ b/table/scanner.go @@ -1644,9 +1644,11 @@ func (scan *Scan) ToArrowRecords(ctx context.Context) (*arrow.Schema, iter.Seq2[ // ReadTasks reads Arrow records from a specific set of FileScanTasks, applying the // scan's projection, per-task residual filters, and delete handling. This is useful // when the caller has already planned or selected specific tasks to read. -// Positional- and equality-delete read errors are returned by the iterator during -// iteration; deletion-vector read errors are returned by ReadTasks before it returns -// an iterator. The returned iterator is single-use. +// Positional- and equality-delete read errors are delivered through the iterator +// only if iteration reaches the task that encounters the error; a row limit or +// early termination may finish the scan before the error is observed. +// Deletion-vector read errors are returned by ReadTasks before it returns an +// iterator. The returned iterator is single-use. func (scan *Scan) ReadTasks(ctx context.Context, tasks []FileScanTask) (*arrow.Schema, iter.Seq2[arrow.RecordBatch, error], error) { if atomic.LoadUint32(&scan.closed) != 0 { return nil, nil, fmt.Errorf("%w: scan is closed", ErrInvalidOperation) From b531f4db0f2706a44eee9b9c1c2be87f2bbc79c2 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Fri, 4 Sep 2026 14:36:05 +0200 Subject: [PATCH 10/10] fix(table): include position delete paths in eager errors --- table/arrow_scanner.go | 2 +- table/arrow_scanner_test.go | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/table/arrow_scanner.go b/table/arrow_scanner.go index ae8780891..406e8d82a 100644 --- a/table/arrow_scanner.go +++ b/table/arrow_scanner.go @@ -114,7 +114,7 @@ func readAllDeleteFiles(ctx context.Context, fs iceio.IO, tasks []FileScanTask, g.Go(func() error { deletes, err := readDeletes(gctx, fs, v) if err != nil { - return err + return fmt.Errorf("read position deletes from %s: %w", v.FilePath(), err) } if deletes == nil { return nil diff --git a/table/arrow_scanner_test.go b/table/arrow_scanner_test.go index 5dbe358ea..eaca85236 100644 --- a/table/arrow_scanner_test.go +++ b/table/arrow_scanner_test.go @@ -383,6 +383,7 @@ func TestReadAllDeleteFilesReturnsPartialDeletesOnError(t *testing.T) { // the good worker closes first, then the bad worker returns its error. deletesPerFile, err := readAllDeleteFiles(ctx, testFS, tasks, 2) require.Error(t, err) + require.ErrorContains(t, err, "read position deletes from "+badDeletePath) require.NotNil(t, deletesPerFile) require.Contains(t, deletesPerFile, dataPath) require.Len(t, deletesPerFile[dataPath], 1)