From bc37781923bf729bdb0a85f3d5100c50998eb0ab Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Tue, 22 Sep 2026 00:09:22 -0400 Subject: [PATCH 01/10] perf(table): run compaction groups concurrently in RewriteDataFiles RewriteDataFiles runs ExecuteCompactionGroup calls under a bounded errgroup when MaxConcurrency exceeds 1, in both the atomic and partial-progress paths. Results apply in group order. The atomic path keeps its existing behavior of leaving outputs on error. Refs #2040 Signed-off-by: KranzL <50032317+KranzL@users.noreply.github.com> --- table/rewrite_data_files.go | 157 +++++++--- table/rewrite_data_files_test.go | 479 +++++++++++++++++++++++++++++++ 2 files changed, 602 insertions(+), 34 deletions(-) diff --git a/table/rewrite_data_files.go b/table/rewrite_data_files.go index a66e90da2..033d20704 100644 --- a/table/rewrite_data_files.go +++ b/table/rewrite_data_files.go @@ -28,6 +28,7 @@ import ( "github.com/apache/iceberg-go" iceberginternal "github.com/apache/iceberg-go/internal" iceio "github.com/apache/iceberg-go/io" + "golang.org/x/sync/errgroup" ) // RewriteResult summarizes a completed compaction. @@ -217,6 +218,22 @@ type RewriteDataFilesOptions struct { // size, scan concurrency). See the With* helpers returning // [CompactionGroupOption]. GroupOptions []CompactionGroupOption + + // MaxConcurrency bounds how many compaction groups run at once. + // Zero and one both mean sequential execution, which is the default. + // Larger values run [ExecuteCompactionGroup] calls under a bounded + // errgroup and apply their results in the original group order, so + // manifests and [RewriteResult] are identical to a sequential run. + // The first error cancels the groups still running and is returned. + // Peak record-pipeline memory is MaxConcurrency times the per-group + // bound stated on [WithCompactionArrowBatchSize]: + // + // MaxConcurrency x (workers x (rows in the largest task + n) + (recordBatchBufferSize + 2) x n) + // + // rows, where workers, n and recordBatchBufferSize are the per-group + // values. Delete-side memory is outside this bound. Negative values + // are rejected with [ErrInvalidOperation]. + MaxConcurrency int } // CompactionGroupOption configures a single [ExecuteCompactionGroup] @@ -337,6 +354,9 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT if _, err := t.txnMeta(); err != nil { return nil, err } + if opts.MaxConcurrency < 0 { + return nil, fmt.Errorf("%w: MaxConcurrency must be non-negative", ErrInvalidOperation) + } if opts.PartialProgress { return t.rewriteDataFilesPartial(ctx, groups, opts) } @@ -348,31 +368,51 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT rewrite := t.NewRewrite(opts.SnapshotProps) stagedDeleteFiles := make(map[string]struct{}) - for _, group := range groups { - if err := ctx.Err(); err != nil { + if opts.MaxConcurrency > 1 { + results, err := executeCompactionGroups(ctx, t.tbl, groups, opts.GroupOptions, opts.MaxConcurrency) + if err != nil { return result, err } - - if len(group.Tasks) == 0 { - continue + for _, gr := range results { + if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { + continue + } + rewrite.ApplyResult(gr) + accumulateGroupMetrics(result, gr) + for _, df := range gr.SafePosDeletes { + stagedDeleteFiles[df.FilePath()] = struct{}{} + } + for _, df := range gr.SafeDeletionVectors { + stagedDeleteFiles[df.FilePath()] = struct{}{} + } } + } else { + for _, group := range groups { + if err := ctx.Err(); err != nil { + return result, err + } - gr, err := ExecuteCompactionGroup(ctx, t.tbl, group, opts.GroupOptions...) - if err != nil { - return result, err - } + if len(group.Tasks) == 0 { + continue + } - if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { - continue - } + gr, err := ExecuteCompactionGroup(ctx, t.tbl, group, opts.GroupOptions...) + if err != nil { + return result, err + } - rewrite.ApplyResult(gr) - accumulateGroupMetrics(result, gr) - for _, df := range gr.SafePosDeletes { - stagedDeleteFiles[df.FilePath()] = struct{}{} - } - for _, df := range gr.SafeDeletionVectors { - stagedDeleteFiles[df.FilePath()] = struct{}{} + if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { + continue + } + + rewrite.ApplyResult(gr) + accumulateGroupMetrics(result, gr) + for _, df := range gr.SafePosDeletes { + stagedDeleteFiles[df.FilePath()] = struct{}{} + } + for _, df := range gr.SafeDeletionVectors { + stagedDeleteFiles[df.FilePath()] = struct{}{} + } } } @@ -407,6 +447,35 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT return result, nil } +func executeCompactionGroups(ctx context.Context, tbl *Table, groups []CompactionTaskGroup, groupOpts []CompactionGroupOption, maxConcurrency int) ([]CompactionGroupResult, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + g, gctx := errgroup.WithContext(ctx) + g.SetLimit(min(maxConcurrency, len(groups))) + results := make([]CompactionGroupResult, len(groups)) + for i, group := range groups { + if len(group.Tasks) == 0 { + continue + } + g.Go(func() error { + gr, err := ExecuteCompactionGroup(gctx, tbl, group, groupOpts...) + results[i] = gr + + return err + }) + } + if err := g.Wait(); err != nil { + if ctx.Err() != nil { + return results, ctx.Err() + } + + return results, err + } + + return results, nil +} + // ExecuteCompactionGroup reads a compaction group's tasks (with // deletes applied), writes consolidated output files via // [WriteRecords], and computes the position-delete files safe to @@ -639,26 +708,46 @@ func (t *Transaction) rewriteDataFilesPartial(ctx context.Context, groups []Comp return cause } - for _, group := range batchGroups { - if err := ctx.Err(); err != nil { - return result, cleanupBatch(err) - } - - gr, err := ExecuteCompactionGroup(ctx, current, group, opts.GroupOptions...) + if opts.MaxConcurrency > 1 { + results, err := executeCompactionGroups(ctx, current, batchGroups, opts.GroupOptions, opts.MaxConcurrency) if err != nil { - return result, cleanupBatch(err, gr) + return result, cleanupBatch(err, results...) } - - if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { - continue + for _, gr := range results { + if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { + continue + } + batchResults = append(batchResults, gr) + for _, df := range gr.OldDataFiles { + if _, ok := rewrittenPaths[df.FilePath()]; ok { + continue + } + rewrittenPaths[df.FilePath()] = struct{}{} + rewrittenFiles = append(rewrittenFiles, df) + } } - batchResults = append(batchResults, gr) - for _, df := range gr.OldDataFiles { - if _, ok := rewrittenPaths[df.FilePath()]; ok { + } else { + for _, group := range batchGroups { + if err := ctx.Err(); err != nil { + return result, cleanupBatch(err) + } + + gr, err := ExecuteCompactionGroup(ctx, current, group, opts.GroupOptions...) + if err != nil { + return result, cleanupBatch(err, gr) + } + + if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { continue } - rewrittenPaths[df.FilePath()] = struct{}{} - rewrittenFiles = append(rewrittenFiles, df) + batchResults = append(batchResults, gr) + for _, df := range gr.OldDataFiles { + if _, ok := rewrittenPaths[df.FilePath()]; ok { + continue + } + rewrittenPaths[df.FilePath()] = struct{}{} + rewrittenFiles = append(rewrittenFiles, df) + } } } diff --git a/table/rewrite_data_files_test.go b/table/rewrite_data_files_test.go index d9bc9dc6d..1e4003f97 100644 --- a/table/rewrite_data_files_test.go +++ b/table/rewrite_data_files_test.go @@ -23,8 +23,11 @@ import ( "fmt" "os" "path/filepath" + "slices" "strings" + "sync" "testing" + "time" "github.com/apache/arrow-go/v18/arrow" "github.com/apache/arrow-go/v18/arrow/array" @@ -1401,3 +1404,479 @@ func appendEqualityDelete(t *testing.T, tbl *table.Table, equalityFieldIDs []int return out } + +func newMaxConcPartitionedTable(t *testing.T, fs iceio.IO) *table.Table { + t.Helper() + + location := filepath.ToSlash(t.TempDir()) + schema := iceberg.NewSchema(0, + iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64, Required: true}, + iceberg.NestedField{ID: 2, Name: "data", Type: iceberg.PrimitiveTypes.String, Required: false}, + ) + spec := iceberg.NewPartitionSpec(iceberg.PartitionField{ + SourceIDs: []int{2}, FieldID: 1000, Transform: iceberg.IdentityTransform{}, Name: "data", + }) + meta, err := table.NewMetadata(schema, &spec, table.UnsortedSortOrder, location, + iceberg.Properties{table.PropertyFormatVersion: "2"}) + require.NoError(t, err) + + cat := &partialProgressCatalog{metadata: meta} + + return table.New( + table.Identifier{"db", "max_conc_test"}, + meta, location+"/metadata/v1.metadata.json", + func(context.Context) (iceio.IO, error) { return fs, nil }, + cat, + ) +} + +func addMaxConcPartitions(t *testing.T, tbl *table.Table, partitions, filesPerPartition, rowsPerFile int) *table.Table { + t.Helper() + + var nextID int64 = 1 + for p := range partitions { + partition := fmt.Sprintf("p%d", p) + for f := range filesPerPartition { + ids := make([]int64, rowsPerFile) + for r := range rowsPerFile { + ids[r] = nextID + nextID++ + } + tbl = addPartitionedRowsOnRef(t, tbl, table.MainBranch, fmt.Sprintf("p%d-%d", p, f), partition, ids...) + } + } + + return tbl +} + +func groupsByPartition(t *testing.T, tbl *table.Table) []table.CompactionTaskGroup { + t.Helper() + + tasks, err := tbl.Scan().PlanFiles(t.Context()) + require.NoError(t, err) + + byPart := make(map[string][]table.FileScanTask) + for _, task := range tasks { + part, ok := task.File.Partition()[1000].(string) + require.True(t, ok) + byPart[part] = append(byPart[part], task) + } + keys := make([]string, 0, len(byPart)) + for k := range byPart { + keys = append(keys, k) + } + slices.Sort(keys) + + groups := make([]table.CompactionTaskGroup, 0, len(keys)) + for _, k := range keys { + var total int64 + for _, task := range byPart[k] { + total += task.File.FileSizeBytes() + } + groups = append(groups, table.CompactionTaskGroup{ + PartitionKey: k, + Tasks: byPart[k], + TotalSizeBytes: total, + }) + } + + return groups +} + +func rowsByPartitionValue(t *testing.T, tbl *table.Table) map[string]int64 { + t.Helper() + + _, itr, err := tbl.Scan().ToArrowRecords(t.Context()) + require.NoError(t, err) + + out := make(map[string]int64) + for rec, err := range itr { + require.NoError(t, err) + idx := rec.Schema().FieldIndices("data") + require.NotEmpty(t, idx) + col, ok := rec.Column(idx[0]).(*array.String) + require.True(t, ok) + for i := range int(rec.NumRows()) { + out[col.Value(i)]++ + } + rec.Release() + } + + return out +} + +func manifestDataPartitions(t *testing.T, tbl *table.Table) []string { + t.Helper() + + snap := tbl.CurrentSnapshot() + require.NotNil(t, snap) + fs, err := tbl.FS(t.Context()) + require.NoError(t, err) + manifests, err := snap.Manifests(fs) + require.NoError(t, err) + + var parts []string + for _, m := range manifests { + for e, err := range m.Entries(fs, false) { + require.NoError(t, err) + if e.Status() == iceberg.EntryStatusDELETED { + continue + } + df := e.DataFile() + if df.ContentType() != iceberg.EntryContentData { + continue + } + part, ok := df.Partition()[1000].(string) + require.True(t, ok) + parts = append(parts, part) + } + } + + return parts +} + +func manifestLiveDataPaths(t *testing.T, tbl *table.Table) []string { + t.Helper() + + snap := tbl.CurrentSnapshot() + require.NotNil(t, snap) + fs, err := tbl.FS(t.Context()) + require.NoError(t, err) + manifests, err := snap.Manifests(fs) + require.NoError(t, err) + + var paths []string + for _, m := range manifests { + for e, err := range m.Entries(fs, false) { + require.NoError(t, err) + if e.Status() == iceberg.EntryStatusDELETED { + continue + } + df := e.DataFile() + if df.ContentType() != iceberg.EntryContentData { + continue + } + paths = append(paths, df.FilePath()) + } + } + + return paths +} + +func TestRewriteDataFiles_MaxConcurrencyMatchesSequential(t *testing.T) { + tblSeq := newMaxConcPartitionedTable(t, iceio.LocalFS{}) + tblSeq = addMaxConcPartitions(t, tblSeq, 8, 2, 5) + tblConc := newMaxConcPartitionedTable(t, iceio.LocalFS{}) + tblConc = addMaxConcPartitions(t, tblConc, 8, 2, 5) + + groupsSeq := groupsByPartition(t, tblSeq) + groupsConc := groupsByPartition(t, tblConc) + require.Len(t, groupsSeq, 8) + require.Len(t, groupsConc, 8) + + txSeq := tblSeq.NewTransaction() + resSeq, err := txSeq.RewriteDataFiles(t.Context(), groupsSeq, table.RewriteDataFilesOptions{}) + require.NoError(t, err) + committedSeq, err := txSeq.Commit(t.Context()) + require.NoError(t, err) + + txConc := tblConc.NewTransaction() + resConc, err := txConc.RewriteDataFiles(t.Context(), groupsConc, table.RewriteDataFilesOptions{MaxConcurrency: 4}) + require.NoError(t, err) + committedConc, err := txConc.Commit(t.Context()) + require.NoError(t, err) + + assert.Equal(t, resSeq.RewrittenGroups, resConc.RewrittenGroups) + assert.Equal(t, resSeq.AddedDataFiles, resConc.AddedDataFiles) + assert.Equal(t, resSeq.RemovedDataFiles, resConc.RemovedDataFiles) + assert.Equal(t, resSeq.RemovedPositionDeleteFiles, resConc.RemovedPositionDeleteFiles) + assert.Equal(t, resSeq.RemovedEqualityDeleteFiles, resConc.RemovedEqualityDeleteFiles) + assert.Equal(t, resSeq.RemovedDeletionVectorFiles, resConc.RemovedDeletionVectorFiles) + assert.Equal(t, resSeq.BytesBefore, resConc.BytesBefore) + assert.Equal(t, 8, resConc.RewrittenGroups) + assert.Equal(t, 16, resConc.RemovedDataFiles) + assert.Equal(t, 8, resConc.AddedDataFiles) + + assert.Equal(t, rowsByPartitionValue(t, committedSeq), rowsByPartitionValue(t, committedConc)) + for p := range 8 { + assert.Equal(t, int64(10), rowsByPartitionValue(t, committedConc)[fmt.Sprintf("p%d", p)]) + } + + paths := manifestLiveDataPaths(t, committedConc) + require.Len(t, paths, 8) + assert.Len(t, map[string]struct{}{paths[0]: {}, paths[1]: {}, paths[2]: {}, paths[3]: {}, paths[4]: {}, paths[5]: {}, paths[6]: {}, paths[7]: {}}, 8) +} + +func TestRewriteDataFiles_MaxConcurrencyNegativeRejected(t *testing.T) { + tbl := newRewriteTestTable(t) + + tx := tbl.NewTransaction() + _, err := tx.RewriteDataFiles(t.Context(), nil, table.RewriteDataFilesOptions{MaxConcurrency: -1}) + require.ErrorIs(t, err, table.ErrInvalidOperation) + + txPartial := tbl.NewTransaction() + _, err = txPartial.RewriteDataFiles(t.Context(), nil, table.RewriteDataFilesOptions{PartialProgress: true, MaxConcurrency: -1}) + require.ErrorIs(t, err, table.ErrInvalidOperation) +} + +func TestRewriteDataFiles_MaxConcurrencyDeterministicOrder(t *testing.T) { + tblA := newMaxConcPartitionedTable(t, iceio.LocalFS{}) + tblA = addMaxConcPartitions(t, tblA, 8, 1, 5) + tblB := newMaxConcPartitionedTable(t, iceio.LocalFS{}) + tblB = addMaxConcPartitions(t, tblB, 8, 1, 5) + + groupsA := groupsByPartition(t, tblA) + groupsB := groupsByPartition(t, tblB) + + txA := tblA.NewTransaction() + _, err := txA.RewriteDataFiles(t.Context(), groupsA, table.RewriteDataFilesOptions{MaxConcurrency: 4}) + require.NoError(t, err) + committedA, err := txA.Commit(t.Context()) + require.NoError(t, err) + + txB := tblB.NewTransaction() + _, err = txB.RewriteDataFiles(t.Context(), groupsB, table.RewriteDataFilesOptions{MaxConcurrency: 4}) + require.NoError(t, err) + committedB, err := txB.Commit(t.Context()) + require.NoError(t, err) + + orderA := manifestDataPartitions(t, committedA) + orderB := manifestDataPartitions(t, committedB) + require.Len(t, orderA, 8) + require.Len(t, orderB, 8) + assert.Equal(t, orderA, orderB) + assert.Equal(t, []string{"p0", "p1", "p2", "p3", "p4", "p5", "p6", "p7"}, orderA) +} + +type failOpenIO struct { + iceio.LocalFS + mu sync.Mutex + failSubstr string + failErr error +} + +func (f *failOpenIO) setFail(substr string, err error) { + f.mu.Lock() + defer f.mu.Unlock() + f.failSubstr = substr + f.failErr = err +} + +func (f *failOpenIO) Open(name string) (iceio.File, error) { + f.mu.Lock() + substr, failErr := f.failSubstr, f.failErr + f.mu.Unlock() + if substr != "" && strings.Contains(name, substr) { + return nil, failErr + } + + return f.LocalFS.Open(name) +} + +func TestRewriteDataFiles_MaxConcurrencyGroupFailure(t *testing.T) { + injected := errors.New("injected compaction read failure") + + fsAtomic := &failOpenIO{} + tblAtomic := newMaxConcPartitionedTable(t, fsAtomic) + tblAtomic = addMaxConcPartitions(t, tblAtomic, 4, 1, 5) + groupsAtomic := groupsByPartition(t, tblAtomic) + require.Len(t, groupsAtomic, 4) + fsAtomic.setFail(groupsAtomic[2].Tasks[0].File.FilePath(), injected) + + txAtomic := tblAtomic.NewTransaction() + _, err := txAtomic.RewriteDataFiles(t.Context(), groupsAtomic, table.RewriteDataFilesOptions{MaxConcurrency: 4}) + require.Error(t, err) + assert.Contains(t, err.Error(), injected.Error()) + + fsPartial := &failOpenIO{} + tblPartial := newMaxConcPartitionedTable(t, fsPartial) + tblPartial = addMaxConcPartitions(t, tblPartial, 4, 1, 5) + groupsPartial := groupsByPartition(t, tblPartial) + require.Len(t, groupsPartial, 4) + fsPartial.setFail(groupsPartial[1].Tasks[0].File.FilePath(), injected) + beforeFiles := parquetFiles(t, tblPartial.Location()) + + txPartial := tblPartial.NewTransaction() + result, err := txPartial.RewriteDataFiles(t.Context(), groupsPartial, table.RewriteDataFilesOptions{ + PartialProgress: true, + MaxCommits: 1, + MaxConcurrency: 4, + }) + require.Error(t, err) + require.NotNil(t, result) + assert.Contains(t, err.Error(), injected.Error()) + assert.Empty(t, result.CompletedGroups) + assert.ElementsMatch(t, beforeFiles, parquetFiles(t, tblPartial.Location())) +} + +type gateOpenIO struct { + iceio.LocalFS + mu sync.Mutex + enabled bool + entered chan struct{} + release chan struct{} +} + +func (g *gateOpenIO) enable() { + g.mu.Lock() + defer g.mu.Unlock() + g.enabled = true + g.entered = make(chan struct{}, 32) + g.release = make(chan struct{}) +} + +func (g *gateOpenIO) releaseAll() { + g.mu.Lock() + defer g.mu.Unlock() + select { + case <-g.release: + default: + close(g.release) + } +} + +func (g *gateOpenIO) Open(name string) (iceio.File, error) { + g.mu.Lock() + enabled, entered, release := g.enabled, g.entered, g.release + g.mu.Unlock() + if enabled && strings.Contains(name, "/data/") { + select { + case <-release: + default: + select { + case entered <- struct{}{}: + default: + } + <-release + } + } + + return g.LocalFS.Open(name) +} + +func TestRewriteDataFiles_MaxConcurrencyContextCancel(t *testing.T) { + for _, partial := range []bool{false, true} { + gate := &gateOpenIO{} + tbl := newMaxConcPartitionedTable(t, gate) + tbl = addMaxConcPartitions(t, tbl, 8, 1, 10) + groups := groupsByPartition(t, tbl) + require.Len(t, groups, 8) + gate.enable() + + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan struct{}) + var rewriteErr error + go func() { + defer close(done) + tx := tbl.NewTransaction() + opts := table.RewriteDataFilesOptions{MaxConcurrency: 4} + if partial { + opts.PartialProgress = true + opts.MaxCommits = 1 + } + _, rewriteErr = tx.RewriteDataFiles(ctx, groups, opts) + }() + + for range 4 { + select { + case <-gate.entered: + case <-done: + t.Fatalf("rewrite finished before 4 groups were in flight, err=%v", rewriteErr) + case <-t.Context().Done(): + t.Fatal("test context done while waiting for groups") + } + } + cancel() + gate.releaseAll() + <-done + require.Error(t, rewriteErr) + assert.ErrorIs(t, rewriteErr, context.Canceled) + assert.Equal(t, ctx.Err(), rewriteErr) + } +} + +type countOpenFile struct { + iceio.File + owner *countOpenIO + once sync.Once +} + +func (f *countOpenFile) Close() error { + err := f.File.Close() + f.once.Do(func() { + f.owner.mu.Lock() + defer f.owner.mu.Unlock() + f.owner.cur-- + }) + + return err +} + +type countOpenIO struct { + iceio.LocalFS + mu sync.Mutex + cur int + peak int +} + +func (c *countOpenIO) Open(name string) (iceio.File, error) { + f, err := c.LocalFS.Open(name) + if err != nil { + return nil, err + } + if strings.Contains(name, "/data/") { + c.mu.Lock() + c.cur++ + if c.cur > c.peak { + c.peak = c.cur + } + c.mu.Unlock() + time.Sleep(20 * time.Millisecond) + + return &countOpenFile{File: f, owner: c}, nil + } + + return f, nil +} + +func (c *countOpenIO) reset() { + c.mu.Lock() + defer c.mu.Unlock() + c.cur = 0 + c.peak = 0 +} + +func (c *countOpenIO) getPeak() int { + c.mu.Lock() + defer c.mu.Unlock() + + return c.peak +} + +func TestRewriteDataFiles_MaxConcurrencyLimitsInFlight(t *testing.T) { + for _, maxConc := range []int{4, 0, 1} { + counter := &countOpenIO{} + tbl := newMaxConcPartitionedTable(t, counter) + tbl = addMaxConcPartitions(t, tbl, 8, 1, 10) + groups := groupsByPartition(t, tbl) + require.Len(t, groups, 8) + counter.reset() + + tx := tbl.NewTransaction() + _, err := tx.RewriteDataFiles(t.Context(), groups, table.RewriteDataFilesOptions{ + MaxConcurrency: maxConc, + GroupOptions: []table.CompactionGroupOption{table.WithCompactionScanConcurrency(1)}, + }) + require.NoError(t, err) + _, err = tx.Commit(t.Context()) + require.NoError(t, err) + + peak := counter.getPeak() + if maxConc > 1 { + assert.LessOrEqual(t, peak, maxConc) + assert.GreaterOrEqual(t, peak, 2) + } else { + assert.Equal(t, 1, peak) + } + } +} From 23ef481e057cbeb4b3b1fb5b23ec61791663640d Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Tue, 22 Sep 2026 00:20:38 -0400 Subject: [PATCH 02/10] perf(table): add multi-group scaling benchmark for concurrent compaction BenchmarkRewriteDataFilesGroupConcurrency in table/rewrite_data_files_bench_test.go rewrites 8 partitions of 8 files over MaxConcurrency 1, 2, 4, 8 and reports ns/op, B/op and rows/s. Each iteration runs on a fresh transaction and removes its outputs, so iterations stay comparable. The MaxConcurrency doc comment gains the byte form of the sizing bound. Refs #2040. Signed-off-by: KranzL <50032317+KranzL@users.noreply.github.com> --- table/rewrite_data_files.go | 3 +- table/rewrite_data_files_bench_test.go | 243 +++++++++++++++++++++++++ 2 files changed, 245 insertions(+), 1 deletion(-) create mode 100644 table/rewrite_data_files_bench_test.go diff --git a/table/rewrite_data_files.go b/table/rewrite_data_files.go index 033d20704..d3cf22dbf 100644 --- a/table/rewrite_data_files.go +++ b/table/rewrite_data_files.go @@ -231,7 +231,8 @@ type RewriteDataFilesOptions struct { // MaxConcurrency x (workers x (rows in the largest task + n) + (recordBatchBufferSize + 2) x n) // // rows, where workers, n and recordBatchBufferSize are the per-group - // values. Delete-side memory is outside this bound. Negative values + // values. Multiply rows by the average row width in bytes for a byte + // estimate. Delete-side memory is outside this bound. Negative values // are rejected with [ErrInvalidOperation]. MaxConcurrency int } diff --git a/table/rewrite_data_files_bench_test.go b/table/rewrite_data_files_bench_test.go new file mode 100644 index 000000000..2c42eedcb --- /dev/null +++ b/table/rewrite_data_files_bench_test.go @@ -0,0 +1,243 @@ +// 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_test + +import ( + "context" + "fmt" + "math/rand/v2" + "os" + "path/filepath" + "slices" + "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/pqarrow" + "github.com/apache/iceberg-go" + iceio "github.com/apache/iceberg-go/io" + "github.com/apache/iceberg-go/table" + "github.com/stretchr/testify/require" +) + +const ( + groupConcPartitions = 8 + groupConcFilesPerPart = 8 + groupConcRowsPerFile = 30000 + groupConcPayloadWords = 6 + groupConcPartitionField = 1000 +) + +func newGroupConcTable(tb testing.TB) *table.Table { + tb.Helper() + + location := filepath.ToSlash(tb.TempDir()) + schema := iceberg.NewSchema(0, + iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64, Required: true}, + iceberg.NestedField{ID: 2, Name: "data", Type: iceberg.PrimitiveTypes.String, Required: false}, + iceberg.NestedField{ID: 3, Name: "payload", Type: iceberg.PrimitiveTypes.String, Required: false}, + iceberg.NestedField{ID: 4, Name: "score", Type: iceberg.PrimitiveTypes.Float64, Required: false}, + ) + spec := iceberg.NewPartitionSpec(iceberg.PartitionField{ + SourceIDs: []int{2}, FieldID: groupConcPartitionField, Transform: iceberg.IdentityTransform{}, Name: "data", + }) + meta, err := table.NewMetadata(schema, &spec, table.UnsortedSortOrder, location, + iceberg.Properties{table.PropertyFormatVersion: "2"}) + require.NoError(tb, err) + + cat := &partialProgressCatalog{metadata: meta} + + return table.New( + table.Identifier{"db", "group_concurrency_bench"}, + meta, location+"/metadata/v1.metadata.json", + func(context.Context) (iceio.IO, error) { return iceio.LocalFS{}, nil }, + cat, + ) +} + +func writeGroupConcFile(tb testing.TB, path string, sc *arrow.Schema, partition string, fileIdx int) int64 { + tb.Helper() + + mem := memory.DefaultAllocator + rng := rand.New(rand.NewPCG(uint64(fileIdx), 0x9e3779b97f4a7c15)) + + idB := array.NewInt64Builder(mem) + dataB := array.NewStringBuilder(mem) + payloadB := array.NewStringBuilder(mem) + scoreB := array.NewFloat64Builder(mem) + defer idB.Release() + defer dataB.Release() + defer payloadB.Release() + defer scoreB.Release() + + idB.Reserve(groupConcRowsPerFile) + dataB.Reserve(groupConcRowsPerFile) + payloadB.Reserve(groupConcRowsPerFile) + scoreB.Reserve(groupConcRowsPerFile) + + var scratch [groupConcPayloadWords]uint64 + for i := range groupConcRowsPerFile { + id := int64(fileIdx*groupConcRowsPerFile + i) + idB.Append(id) + dataB.Append(partition) + for w := range groupConcPayloadWords { + scratch[w] = rng.Uint64() + } + payloadB.Append(fmt.Sprintf("%016x%016x%016x%016x%016x%016x", + scratch[0], scratch[1], scratch[2], scratch[3], scratch[4], scratch[5])) + scoreB.Append(float64(id) * 0.5) + } + + rec := array.NewRecordBatch(sc, []arrow.Array{ + idB.NewArray(), dataB.NewArray(), payloadB.NewArray(), scoreB.NewArray(), + }, int64(groupConcRowsPerFile)) + defer rec.Release() + + fs := iceio.LocalFS{} + fw, err := fs.Create(path) + require.NoError(tb, err) + defer fw.Close() + + arrTable := array.NewTableFromRecords(sc, []arrow.RecordBatch{rec}) + defer arrTable.Release() + + props := parquet.NewWriterProperties(parquet.WithStats(true)) + require.NoError(tb, pqarrow.WriteTable(arrTable, fw, int64(groupConcRowsPerFile), props, pqarrow.DefaultWriterProps())) + + info, err := os.Stat(path) + require.NoError(tb, err) + + return info.Size() +} + +func planGroupConcGroups(tb testing.TB, tbl *table.Table) []table.CompactionTaskGroup { + tb.Helper() + + tasks, err := tbl.Scan().PlanFiles(context.Background()) + require.NoError(tb, err) + require.Len(tb, tasks, groupConcPartitions*groupConcFilesPerPart) + + byPart := make(map[string][]table.FileScanTask, groupConcPartitions) + for _, task := range tasks { + part, ok := task.File.Partition()[groupConcPartitionField].(string) + require.True(tb, ok) + byPart[part] = append(byPart[part], task) + } + + keys := make([]string, 0, len(byPart)) + for k := range byPart { + keys = append(keys, k) + } + slices.Sort(keys) + require.Len(tb, keys, groupConcPartitions) + + groups := make([]table.CompactionTaskGroup, 0, len(keys)) + for _, k := range keys { + require.Len(tb, byPart[k], groupConcFilesPerPart) + var total int64 + for _, task := range byPart[k] { + total += task.File.FileSizeBytes() + } + groups = append(groups, table.CompactionTaskGroup{ + PartitionKey: k, + Tasks: byPart[k], + TotalSizeBytes: total, + }) + } + + return groups +} + +func groupConcOutputPaths(tb testing.TB, location string, before map[string]struct{}) []string { + tb.Helper() + + paths, err := filepath.Glob(filepath.Join(location, "data", "*.parquet")) + require.NoError(tb, err) + + var out []string + for _, p := range paths { + if _, ok := before[p]; !ok { + out = append(out, p) + } + } + + return out +} + +func BenchmarkRewriteDataFilesGroupConcurrency(b *testing.B) { + tbl := newGroupConcTable(b) + + arrowSc, err := table.SchemaToArrowSchema(tbl.Schema(), nil, false, false) + require.NoError(b, err) + + ctx := context.Background() + for p := range groupConcPartitions { + partition := fmt.Sprintf("p%d", p) + files := make([]iceberg.DataFile, 0, groupConcFilesPerPart) + for f := range groupConcFilesPerPart { + fileIdx := p*groupConcFilesPerPart + f + dataPath := tbl.Location() + "/data/" + fmt.Sprintf("p%d-file-%d.parquet", p, f) + size := writeGroupConcFile(b, dataPath, arrowSc, partition, fileIdx) + builder, err := iceberg.NewDataFileBuilder( + tbl.Spec(), iceberg.EntryContentData, dataPath, iceberg.ParquetFile, + map[int]any{groupConcPartitionField: partition}, nil, nil, groupConcRowsPerFile, size) + require.NoError(b, err) + files = append(files, builder.Build()) + } + txn := tbl.NewTransaction() + require.NoError(b, txn.AddDataFiles(ctx, files, nil)) + tbl, err = txn.Commit(ctx) + require.NoError(b, err) + } + + groups := planGroupConcGroups(b, tbl) + totalRows := int64(groupConcPartitions * groupConcFilesPerPart * groupConcRowsPerFile) + + for _, maxConcurrency := range []int{1, 2, 4, 8} { + b.Run(fmt.Sprintf("MaxConcurrency=%d", maxConcurrency), func(b *testing.B) { + before, err := filepath.Glob(filepath.Join(tbl.Location(), "data", "*.parquet")) + require.NoError(b, err) + inputs := make(map[string]struct{}, len(before)) + for _, p := range before { + inputs[p] = struct{}{} + } + opts := table.RewriteDataFilesOptions{MaxConcurrency: maxConcurrency} + + b.ReportAllocs() + b.ResetTimer() + var completed int64 + for b.Loop() { + tx := tbl.NewTransaction() + result, err := tx.RewriteDataFiles(ctx, groups, opts) + require.NoError(b, err) + require.Equal(b, groupConcPartitions, result.RewrittenGroups) + completed++ + + b.StopTimer() + for _, p := range groupConcOutputPaths(b, tbl.Location(), inputs) { + require.NoError(b, os.Remove(p)) + } + b.StartTimer() + } + b.StopTimer() + b.ReportMetric(float64(totalRows*completed)/b.Elapsed().Seconds(), "rows/s") + }) + } +} From e6253826b148ac3c2f18ddae5330721f6bb65652 Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Tue, 22 Sep 2026 00:54:31 -0400 Subject: [PATCH 03/10] test(table): check concurrent rewrite cleanup with a recursive parquet walk parquetFiles globs location/data/*.parquet and never sees partitioned rewrite outputs under data/data=/. The partial-progress half of the group failure test now uses a recursive WalkDir helper for its before and after file sets, and the sequential parity test asserts the walk contains every live manifest path. Refs #2040 Signed-off-by: KranzL <50032317+KranzL@users.noreply.github.com> --- table/rewrite_data_files_test.go | 28 ++++++++++++++++++++++++++-- 1 file changed, 26 insertions(+), 2 deletions(-) diff --git a/table/rewrite_data_files_test.go b/table/rewrite_data_files_test.go index 1e4003f97..884a8ebf7 100644 --- a/table/rewrite_data_files_test.go +++ b/table/rewrite_data_files_test.go @@ -21,6 +21,7 @@ import ( "context" "errors" "fmt" + "io/fs" "os" "path/filepath" "slices" @@ -846,6 +847,25 @@ func parquetFiles(t *testing.T, location string) []string { return paths } +func allParquetFiles(t *testing.T, location string) []string { + t.Helper() + + var paths []string + err := filepath.WalkDir(filepath.Join(location, "data"), func(path string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + if !d.IsDir() && strings.HasSuffix(path, ".parquet") { + paths = append(paths, filepath.ToSlash(path)) + } + + return nil + }) + require.NoError(t, err) + + return paths +} + func newPartialProgressPartitionedTable(t *testing.T) *table.Table { t.Helper() @@ -1605,6 +1625,10 @@ func TestRewriteDataFiles_MaxConcurrencyMatchesSequential(t *testing.T) { paths := manifestLiveDataPaths(t, committedConc) require.Len(t, paths, 8) assert.Len(t, map[string]struct{}{paths[0]: {}, paths[1]: {}, paths[2]: {}, paths[3]: {}, paths[4]: {}, paths[5]: {}, paths[6]: {}, paths[7]: {}}, 8) + onDisk := allParquetFiles(t, committedConc.Location()) + for _, p := range paths { + assert.Contains(t, onDisk, p) + } } func TestRewriteDataFiles_MaxConcurrencyNegativeRejected(t *testing.T) { @@ -1694,7 +1718,7 @@ func TestRewriteDataFiles_MaxConcurrencyGroupFailure(t *testing.T) { groupsPartial := groupsByPartition(t, tblPartial) require.Len(t, groupsPartial, 4) fsPartial.setFail(groupsPartial[1].Tasks[0].File.FilePath(), injected) - beforeFiles := parquetFiles(t, tblPartial.Location()) + beforeFiles := allParquetFiles(t, tblPartial.Location()) txPartial := tblPartial.NewTransaction() result, err := txPartial.RewriteDataFiles(t.Context(), groupsPartial, table.RewriteDataFilesOptions{ @@ -1706,7 +1730,7 @@ func TestRewriteDataFiles_MaxConcurrencyGroupFailure(t *testing.T) { require.NotNil(t, result) assert.Contains(t, err.Error(), injected.Error()) assert.Empty(t, result.CompletedGroups) - assert.ElementsMatch(t, beforeFiles, parquetFiles(t, tblPartial.Location())) + assert.ElementsMatch(t, beforeFiles, allParquetFiles(t, tblPartial.Location())) } type gateOpenIO struct { From 18958a9238150b375142c1c648f48d30536bed11 Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Tue, 22 Sep 2026 00:56:10 -0400 Subject: [PATCH 04/10] fix(table): collect partitioned outputs when cleaning the group benchmark The clustered writer puts this partitioned table's outputs under data/data=pN/, so the flat data/*.parquet glob in groupConcOutputPaths found nothing and each iteration leaked about 223 MB into the temp dir. Walk location/data recursively and collect every .parquet path outside the input set. Removal stays inside the StopTimer/StartTimer window. Refs #2040. Signed-off-by: KranzL <50032317+KranzL@users.noreply.github.com> --- table/rewrite_data_files_bench_test.go | 21 ++++++++++++++------- 1 file changed, 14 insertions(+), 7 deletions(-) diff --git a/table/rewrite_data_files_bench_test.go b/table/rewrite_data_files_bench_test.go index 2c42eedcb..5037ad46a 100644 --- a/table/rewrite_data_files_bench_test.go +++ b/table/rewrite_data_files_bench_test.go @@ -20,6 +20,7 @@ package table_test import ( "context" "fmt" + "io/fs" "math/rand/v2" "os" "path/filepath" @@ -168,15 +169,21 @@ func planGroupConcGroups(tb testing.TB, tbl *table.Table) []table.CompactionTask func groupConcOutputPaths(tb testing.TB, location string, before map[string]struct{}) []string { tb.Helper() - paths, err := filepath.Glob(filepath.Join(location, "data", "*.parquet")) - require.NoError(tb, err) - var out []string - for _, p := range paths { - if _, ok := before[p]; !ok { - out = append(out, p) + err := filepath.WalkDir(filepath.Join(location, "data"), func(path string, entry fs.DirEntry, err error) error { + if err != nil { + return err } - } + if entry.IsDir() || filepath.Ext(path) != ".parquet" { + return nil + } + if _, ok := before[path]; !ok { + out = append(out, path) + } + + return nil + }) + require.NoError(tb, err) return out } From 2b706fbc0af3e26288ae2d194b79f024ccaad9a1 Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Wed, 23 Sep 2026 17:38:33 -0400 Subject: [PATCH 05/10] perf(table): clean up atomic rewrite outputs on group failure Remove outputs the atomic path wrote before a group failed, in both the concurrent and the sequential branch, mirroring the partial-progress cleanup. Factor the per-group apply and batch accumulation blocks into applyAtomicGroupResult and appendPartialGroupResult so the four executor branches share one copy. Only substitute the caller context error in executeCompactionGroups when the group error is itself context-caused, and clamp the errgroup limit to at least one. Signed-off-by: KranzL <50032317+KranzL@users.noreply.github.com> --- table/rewrite_data_files.go | 109 +++++++++++++++++++----------------- 1 file changed, 58 insertions(+), 51 deletions(-) diff --git a/table/rewrite_data_files.go b/table/rewrite_data_files.go index d3cf22dbf..6dd1c70d7 100644 --- a/table/rewrite_data_files.go +++ b/table/rewrite_data_files.go @@ -372,25 +372,16 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT if opts.MaxConcurrency > 1 { results, err := executeCompactionGroups(ctx, t.tbl, groups, opts.GroupOptions, opts.MaxConcurrency) if err != nil { - return result, err + return result, cleanupAtomicRewriteOutputs(ctx, t.tbl, results, err) } for _, gr := range results { - if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { - continue - } - rewrite.ApplyResult(gr) - accumulateGroupMetrics(result, gr) - for _, df := range gr.SafePosDeletes { - stagedDeleteFiles[df.FilePath()] = struct{}{} - } - for _, df := range gr.SafeDeletionVectors { - stagedDeleteFiles[df.FilePath()] = struct{}{} - } + applyAtomicGroupResult(rewrite, result, stagedDeleteFiles, gr) } } else { + var applied []CompactionGroupResult for _, group := range groups { if err := ctx.Err(); err != nil { - return result, err + return result, cleanupAtomicRewriteOutputs(ctx, t.tbl, applied, err) } if len(group.Tasks) == 0 { @@ -399,21 +390,11 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT gr, err := ExecuteCompactionGroup(ctx, t.tbl, group, opts.GroupOptions...) if err != nil { - return result, err - } - - if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { - continue + return result, cleanupAtomicRewriteOutputs(ctx, t.tbl, append(applied, gr), err) } + applied = append(applied, gr) - rewrite.ApplyResult(gr) - accumulateGroupMetrics(result, gr) - for _, df := range gr.SafePosDeletes { - stagedDeleteFiles[df.FilePath()] = struct{}{} - } - for _, df := range gr.SafeDeletionVectors { - stagedDeleteFiles[df.FilePath()] = struct{}{} - } + applyAtomicGroupResult(rewrite, result, stagedDeleteFiles, gr) } } @@ -448,12 +429,42 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT return result, nil } +func applyAtomicGroupResult(rewrite *RewriteFiles, result *RewriteResult, stagedDeleteFiles map[string]struct{}, gr CompactionGroupResult) { + if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { + return + } + rewrite.ApplyResult(gr) + accumulateGroupMetrics(result, gr) + for _, df := range gr.SafePosDeletes { + stagedDeleteFiles[df.FilePath()] = struct{}{} + } + for _, df := range gr.SafeDeletionVectors { + stagedDeleteFiles[df.FilePath()] = struct{}{} + } +} + +func cleanupAtomicRewriteOutputs(ctx context.Context, tbl *Table, results []CompactionGroupResult, cause error) error { + fs, err := tbl.fsF(ctx) + if err != nil { + return errors.Join(cause, fmt.Errorf("open table IO to clean up atomic rewrite outputs: %w", err)) + } + if err := cleanupCompactionOutputs(fs, results); err != nil { + return errors.Join(cause, fmt.Errorf("clean up atomic rewrite outputs: %w", err)) + } + + return cause +} + func executeCompactionGroups(ctx context.Context, tbl *Table, groups []CompactionTaskGroup, groupOpts []CompactionGroupOption, maxConcurrency int) ([]CompactionGroupResult, error) { if err := ctx.Err(); err != nil { return nil, err } + limit := min(maxConcurrency, len(groups)) + if limit < 1 { + limit = 1 + } g, gctx := errgroup.WithContext(ctx) - g.SetLimit(min(maxConcurrency, len(groups))) + g.SetLimit(limit) results := make([]CompactionGroupResult, len(groups)) for i, group := range groups { if len(group.Tasks) == 0 { @@ -467,8 +478,8 @@ func executeCompactionGroups(ctx context.Context, tbl *Table, groups []Compactio }) } if err := g.Wait(); err != nil { - if ctx.Err() != nil { - return results, ctx.Err() + if ctxErr := ctx.Err(); ctxErr != nil && (errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded)) { + return results, ctxErr } return results, err @@ -715,17 +726,7 @@ func (t *Transaction) rewriteDataFilesPartial(ctx context.Context, groups []Comp return result, cleanupBatch(err, results...) } for _, gr := range results { - if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { - continue - } - batchResults = append(batchResults, gr) - for _, df := range gr.OldDataFiles { - if _, ok := rewrittenPaths[df.FilePath()]; ok { - continue - } - rewrittenPaths[df.FilePath()] = struct{}{} - rewrittenFiles = append(rewrittenFiles, df) - } + batchResults, rewrittenFiles = appendPartialGroupResult(batchResults, rewrittenPaths, rewrittenFiles, gr) } } else { for _, group := range batchGroups { @@ -738,17 +739,7 @@ func (t *Transaction) rewriteDataFilesPartial(ctx context.Context, groups []Comp return result, cleanupBatch(err, gr) } - if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { - continue - } - batchResults = append(batchResults, gr) - for _, df := range gr.OldDataFiles { - if _, ok := rewrittenPaths[df.FilePath()]; ok { - continue - } - rewrittenPaths[df.FilePath()] = struct{}{} - rewrittenFiles = append(rewrittenFiles, df) - } + batchResults, rewrittenFiles = appendPartialGroupResult(batchResults, rewrittenPaths, rewrittenFiles, gr) } } @@ -842,6 +833,22 @@ func (t *Transaction) rewriteDataFilesPartial(ctx context.Context, groups []Comp return result, nil } +func appendPartialGroupResult(batchResults []CompactionGroupResult, rewrittenPaths map[string]struct{}, rewrittenFiles []iceberg.DataFile, gr CompactionGroupResult) ([]CompactionGroupResult, []iceberg.DataFile) { + if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 { + return batchResults, rewrittenFiles + } + batchResults = append(batchResults, gr) + for _, df := range gr.OldDataFiles { + if _, ok := rewrittenPaths[df.FilePath()]; ok { + continue + } + rewrittenPaths[df.FilePath()] = struct{}{} + rewrittenFiles = append(rewrittenFiles, df) + } + + return batchResults, rewrittenFiles +} + func recordCommittedRewriteBatch( result *RewriteResult, next *Table, From 7657da62bd27b6bbc407c81d0d5f0b111c37e229 Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Wed, 23 Sep 2026 17:49:13 -0400 Subject: [PATCH 06/10] perf(table): rename MaxConcurrency option to MaxConcurrentGroups The old name collides with the per-group scan worker knob behind WithCompactionScanConcurrency. Cross-reference the two knobs in each other's doc and note that file-open fan-out is about N groups times the scan worker count. Signed-off-by: KranzL <50032317+KranzL@users.noreply.github.com> --- table/rewrite_data_files.go | 39 +++++++++++++++----------- table/rewrite_data_files_bench_test.go | 6 ++-- table/rewrite_data_files_test.go | 36 ++++++++++++------------ 3 files changed, 44 insertions(+), 37 deletions(-) diff --git a/table/rewrite_data_files.go b/table/rewrite_data_files.go index 6dd1c70d7..174b02a82 100644 --- a/table/rewrite_data_files.go +++ b/table/rewrite_data_files.go @@ -219,22 +219,26 @@ type RewriteDataFilesOptions struct { // [CompactionGroupOption]. GroupOptions []CompactionGroupOption - // MaxConcurrency bounds how many compaction groups run at once. + // MaxConcurrentGroups bounds how many compaction groups run at once. // Zero and one both mean sequential execution, which is the default. // Larger values run [ExecuteCompactionGroup] calls under a bounded // errgroup and apply their results in the original group order, so // manifests and [RewriteResult] are identical to a sequential run. // The first error cancels the groups still running and is returned. - // Peak record-pipeline memory is MaxConcurrency times the per-group - // bound stated on [WithCompactionArrowBatchSize]: + // Peak record-pipeline memory is MaxConcurrentGroups times the + // per-group bound stated on [WithCompactionArrowBatchSize]: // - // MaxConcurrency x (workers x (rows in the largest task + n) + (recordBatchBufferSize + 2) x n) + // MaxConcurrentGroups x (workers x (rows in the largest task + n) + (recordBatchBufferSize + 2) x n) // // rows, where workers, n and recordBatchBufferSize are the per-group // values. Multiply rows by the average row width in bytes for a byte - // estimate. Delete-side memory is outside this bound. Negative values - // are rejected with [ErrInvalidOperation]. - MaxConcurrency int + // estimate. Delete-side memory is outside this bound. File-open + // fan-out multiplies too: every group scans with up to + // [WithCompactionScanConcurrency] workers, so N groups open about N + // times the scan worker count in files at once. Size the two knobs + // together against connection and file descriptor limits. Negative + // values are rejected with [ErrInvalidOperation]. + MaxConcurrentGroups int } // CompactionGroupOption configures a single [ExecuteCompactionGroup] @@ -270,7 +274,10 @@ func WithCompactionTargetFileSize(size int64) CompactionGroupOption { // The scan runs min(n, number of tasks) workers, and each worker holds // the decoded batches of at most one task until the writer has taken // them, so the worker count multiplies the read-side term of the -// memory bound stated on [WithCompactionArrowBatchSize]. +// memory bound stated on [WithCompactionArrowBatchSize]. It also +// multiplies [RewriteDataFilesOptions.MaxConcurrentGroups] for file +// opens: N groups scan with up to n workers each, so about N times n +// files are open at once. func WithCompactionScanConcurrency(n int) CompactionGroupOption { return func(c *compactionGroupConfig) { c.scanConcurrency = n @@ -355,8 +362,8 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT if _, err := t.txnMeta(); err != nil { return nil, err } - if opts.MaxConcurrency < 0 { - return nil, fmt.Errorf("%w: MaxConcurrency must be non-negative", ErrInvalidOperation) + if opts.MaxConcurrentGroups < 0 { + return nil, fmt.Errorf("%w: MaxConcurrentGroups must be non-negative", ErrInvalidOperation) } if opts.PartialProgress { return t.rewriteDataFilesPartial(ctx, groups, opts) @@ -369,8 +376,8 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT rewrite := t.NewRewrite(opts.SnapshotProps) stagedDeleteFiles := make(map[string]struct{}) - if opts.MaxConcurrency > 1 { - results, err := executeCompactionGroups(ctx, t.tbl, groups, opts.GroupOptions, opts.MaxConcurrency) + if opts.MaxConcurrentGroups > 1 { + results, err := executeCompactionGroups(ctx, t.tbl, groups, opts.GroupOptions, opts.MaxConcurrentGroups) if err != nil { return result, cleanupAtomicRewriteOutputs(ctx, t.tbl, results, err) } @@ -455,11 +462,11 @@ func cleanupAtomicRewriteOutputs(ctx context.Context, tbl *Table, results []Comp return cause } -func executeCompactionGroups(ctx context.Context, tbl *Table, groups []CompactionTaskGroup, groupOpts []CompactionGroupOption, maxConcurrency int) ([]CompactionGroupResult, error) { +func executeCompactionGroups(ctx context.Context, tbl *Table, groups []CompactionTaskGroup, groupOpts []CompactionGroupOption, maxConcurrentGroups int) ([]CompactionGroupResult, error) { if err := ctx.Err(); err != nil { return nil, err } - limit := min(maxConcurrency, len(groups)) + limit := min(maxConcurrentGroups, len(groups)) if limit < 1 { limit = 1 } @@ -720,8 +727,8 @@ func (t *Transaction) rewriteDataFilesPartial(ctx context.Context, groups []Comp return cause } - if opts.MaxConcurrency > 1 { - results, err := executeCompactionGroups(ctx, current, batchGroups, opts.GroupOptions, opts.MaxConcurrency) + if opts.MaxConcurrentGroups > 1 { + results, err := executeCompactionGroups(ctx, current, batchGroups, opts.GroupOptions, opts.MaxConcurrentGroups) if err != nil { return result, cleanupBatch(err, results...) } diff --git a/table/rewrite_data_files_bench_test.go b/table/rewrite_data_files_bench_test.go index 5037ad46a..99dd8ae88 100644 --- a/table/rewrite_data_files_bench_test.go +++ b/table/rewrite_data_files_bench_test.go @@ -217,15 +217,15 @@ func BenchmarkRewriteDataFilesGroupConcurrency(b *testing.B) { groups := planGroupConcGroups(b, tbl) totalRows := int64(groupConcPartitions * groupConcFilesPerPart * groupConcRowsPerFile) - for _, maxConcurrency := range []int{1, 2, 4, 8} { - b.Run(fmt.Sprintf("MaxConcurrency=%d", maxConcurrency), func(b *testing.B) { + for _, maxConcurrentGroups := range []int{1, 2, 4, 8} { + b.Run(fmt.Sprintf("MaxConcurrentGroups=%d", maxConcurrentGroups), func(b *testing.B) { before, err := filepath.Glob(filepath.Join(tbl.Location(), "data", "*.parquet")) require.NoError(b, err) inputs := make(map[string]struct{}, len(before)) for _, p := range before { inputs[p] = struct{}{} } - opts := table.RewriteDataFilesOptions{MaxConcurrency: maxConcurrency} + opts := table.RewriteDataFilesOptions{MaxConcurrentGroups: maxConcurrentGroups} b.ReportAllocs() b.ResetTimer() diff --git a/table/rewrite_data_files_test.go b/table/rewrite_data_files_test.go index 884a8ebf7..ade78f1a1 100644 --- a/table/rewrite_data_files_test.go +++ b/table/rewrite_data_files_test.go @@ -1583,7 +1583,7 @@ func manifestLiveDataPaths(t *testing.T, tbl *table.Table) []string { return paths } -func TestRewriteDataFiles_MaxConcurrencyMatchesSequential(t *testing.T) { +func TestRewriteDataFiles_MaxConcurrentGroupsMatchesSequential(t *testing.T) { tblSeq := newMaxConcPartitionedTable(t, iceio.LocalFS{}) tblSeq = addMaxConcPartitions(t, tblSeq, 8, 2, 5) tblConc := newMaxConcPartitionedTable(t, iceio.LocalFS{}) @@ -1601,7 +1601,7 @@ func TestRewriteDataFiles_MaxConcurrencyMatchesSequential(t *testing.T) { require.NoError(t, err) txConc := tblConc.NewTransaction() - resConc, err := txConc.RewriteDataFiles(t.Context(), groupsConc, table.RewriteDataFilesOptions{MaxConcurrency: 4}) + resConc, err := txConc.RewriteDataFiles(t.Context(), groupsConc, table.RewriteDataFilesOptions{MaxConcurrentGroups: 4}) require.NoError(t, err) committedConc, err := txConc.Commit(t.Context()) require.NoError(t, err) @@ -1631,19 +1631,19 @@ func TestRewriteDataFiles_MaxConcurrencyMatchesSequential(t *testing.T) { } } -func TestRewriteDataFiles_MaxConcurrencyNegativeRejected(t *testing.T) { +func TestRewriteDataFiles_MaxConcurrentGroupsNegativeRejected(t *testing.T) { tbl := newRewriteTestTable(t) tx := tbl.NewTransaction() - _, err := tx.RewriteDataFiles(t.Context(), nil, table.RewriteDataFilesOptions{MaxConcurrency: -1}) + _, err := tx.RewriteDataFiles(t.Context(), nil, table.RewriteDataFilesOptions{MaxConcurrentGroups: -1}) require.ErrorIs(t, err, table.ErrInvalidOperation) txPartial := tbl.NewTransaction() - _, err = txPartial.RewriteDataFiles(t.Context(), nil, table.RewriteDataFilesOptions{PartialProgress: true, MaxConcurrency: -1}) + _, err = txPartial.RewriteDataFiles(t.Context(), nil, table.RewriteDataFilesOptions{PartialProgress: true, MaxConcurrentGroups: -1}) require.ErrorIs(t, err, table.ErrInvalidOperation) } -func TestRewriteDataFiles_MaxConcurrencyDeterministicOrder(t *testing.T) { +func TestRewriteDataFiles_MaxConcurrentGroupsDeterministicOrder(t *testing.T) { tblA := newMaxConcPartitionedTable(t, iceio.LocalFS{}) tblA = addMaxConcPartitions(t, tblA, 8, 1, 5) tblB := newMaxConcPartitionedTable(t, iceio.LocalFS{}) @@ -1653,13 +1653,13 @@ func TestRewriteDataFiles_MaxConcurrencyDeterministicOrder(t *testing.T) { groupsB := groupsByPartition(t, tblB) txA := tblA.NewTransaction() - _, err := txA.RewriteDataFiles(t.Context(), groupsA, table.RewriteDataFilesOptions{MaxConcurrency: 4}) + _, err := txA.RewriteDataFiles(t.Context(), groupsA, table.RewriteDataFilesOptions{MaxConcurrentGroups: 4}) require.NoError(t, err) committedA, err := txA.Commit(t.Context()) require.NoError(t, err) txB := tblB.NewTransaction() - _, err = txB.RewriteDataFiles(t.Context(), groupsB, table.RewriteDataFilesOptions{MaxConcurrency: 4}) + _, err = txB.RewriteDataFiles(t.Context(), groupsB, table.RewriteDataFilesOptions{MaxConcurrentGroups: 4}) require.NoError(t, err) committedB, err := txB.Commit(t.Context()) require.NoError(t, err) @@ -1697,7 +1697,7 @@ func (f *failOpenIO) Open(name string) (iceio.File, error) { return f.LocalFS.Open(name) } -func TestRewriteDataFiles_MaxConcurrencyGroupFailure(t *testing.T) { +func TestRewriteDataFiles_MaxConcurrentGroupsGroupFailure(t *testing.T) { injected := errors.New("injected compaction read failure") fsAtomic := &failOpenIO{} @@ -1708,7 +1708,7 @@ func TestRewriteDataFiles_MaxConcurrencyGroupFailure(t *testing.T) { fsAtomic.setFail(groupsAtomic[2].Tasks[0].File.FilePath(), injected) txAtomic := tblAtomic.NewTransaction() - _, err := txAtomic.RewriteDataFiles(t.Context(), groupsAtomic, table.RewriteDataFilesOptions{MaxConcurrency: 4}) + _, err := txAtomic.RewriteDataFiles(t.Context(), groupsAtomic, table.RewriteDataFilesOptions{MaxConcurrentGroups: 4}) require.Error(t, err) assert.Contains(t, err.Error(), injected.Error()) @@ -1722,9 +1722,9 @@ func TestRewriteDataFiles_MaxConcurrencyGroupFailure(t *testing.T) { txPartial := tblPartial.NewTransaction() result, err := txPartial.RewriteDataFiles(t.Context(), groupsPartial, table.RewriteDataFilesOptions{ - PartialProgress: true, - MaxCommits: 1, - MaxConcurrency: 4, + PartialProgress: true, + MaxCommits: 1, + MaxConcurrentGroups: 4, }) require.Error(t, err) require.NotNil(t, result) @@ -1778,7 +1778,7 @@ func (g *gateOpenIO) Open(name string) (iceio.File, error) { return g.LocalFS.Open(name) } -func TestRewriteDataFiles_MaxConcurrencyContextCancel(t *testing.T) { +func TestRewriteDataFiles_MaxConcurrentGroupsContextCancel(t *testing.T) { for _, partial := range []bool{false, true} { gate := &gateOpenIO{} tbl := newMaxConcPartitionedTable(t, gate) @@ -1793,7 +1793,7 @@ func TestRewriteDataFiles_MaxConcurrencyContextCancel(t *testing.T) { go func() { defer close(done) tx := tbl.NewTransaction() - opts := table.RewriteDataFilesOptions{MaxConcurrency: 4} + opts := table.RewriteDataFilesOptions{MaxConcurrentGroups: 4} if partial { opts.PartialProgress = true opts.MaxCommits = 1 @@ -1877,7 +1877,7 @@ func (c *countOpenIO) getPeak() int { return c.peak } -func TestRewriteDataFiles_MaxConcurrencyLimitsInFlight(t *testing.T) { +func TestRewriteDataFiles_MaxConcurrentGroupsLimitsInFlight(t *testing.T) { for _, maxConc := range []int{4, 0, 1} { counter := &countOpenIO{} tbl := newMaxConcPartitionedTable(t, counter) @@ -1888,8 +1888,8 @@ func TestRewriteDataFiles_MaxConcurrencyLimitsInFlight(t *testing.T) { tx := tbl.NewTransaction() _, err := tx.RewriteDataFiles(t.Context(), groups, table.RewriteDataFilesOptions{ - MaxConcurrency: maxConc, - GroupOptions: []table.CompactionGroupOption{table.WithCompactionScanConcurrency(1)}, + MaxConcurrentGroups: maxConc, + GroupOptions: []table.CompactionGroupOption{table.WithCompactionScanConcurrency(1)}, }) require.NoError(t, err) _, err = tx.Commit(t.Context()) From bae7b26f9ab81723c0a6f1c17a0279f83940c8d2 Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Wed, 23 Sep 2026 17:54:35 -0400 Subject: [PATCH 07/10] test(table): harden concurrent rewrite failure and accounting tests Assert the atomic failure path leaves no parquet behind, compare per-partition ID sets instead of row counts in the sequential-equivalence test, replace the 20ms sleep in the in-flight test with a deterministic two-open barrier, and run the table cases as named subtests. Signed-off-by: KranzL <50032317+KranzL@users.noreply.github.com> --- table/rewrite_data_files_test.go | 174 +++++++++++++++++++------------ 1 file changed, 105 insertions(+), 69 deletions(-) diff --git a/table/rewrite_data_files_test.go b/table/rewrite_data_files_test.go index ade78f1a1..a5b747979 100644 --- a/table/rewrite_data_files_test.go +++ b/table/rewrite_data_files_test.go @@ -28,7 +28,6 @@ import ( "strings" "sync" "testing" - "time" "github.com/apache/arrow-go/v18/arrow" "github.com/apache/arrow-go/v18/arrow/array" @@ -1503,21 +1502,25 @@ func groupsByPartition(t *testing.T, tbl *table.Table) []table.CompactionTaskGro return groups } -func rowsByPartitionValue(t *testing.T, tbl *table.Table) map[string]int64 { +func idsByPartitionValue(t *testing.T, tbl *table.Table) map[string][]int64 { t.Helper() _, itr, err := tbl.Scan().ToArrowRecords(t.Context()) require.NoError(t, err) - out := make(map[string]int64) + out := make(map[string][]int64) for rec, err := range itr { require.NoError(t, err) - idx := rec.Schema().FieldIndices("data") - require.NotEmpty(t, idx) - col, ok := rec.Column(idx[0]).(*array.String) + dataIdx := rec.Schema().FieldIndices("data") + require.NotEmpty(t, dataIdx) + idIdx := rec.Schema().FieldIndices("id") + require.NotEmpty(t, idIdx) + dataCol, ok := rec.Column(dataIdx[0]).(*array.String) + require.True(t, ok) + idCol, ok := rec.Column(idIdx[0]).(*array.Int64) require.True(t, ok) for i := range int(rec.NumRows()) { - out[col.Value(i)]++ + out[dataCol.Value(i)] = append(out[dataCol.Value(i)], idCol.Value(i)) } rec.Release() } @@ -1617,9 +1620,13 @@ func TestRewriteDataFiles_MaxConcurrentGroupsMatchesSequential(t *testing.T) { assert.Equal(t, 16, resConc.RemovedDataFiles) assert.Equal(t, 8, resConc.AddedDataFiles) - assert.Equal(t, rowsByPartitionValue(t, committedSeq), rowsByPartitionValue(t, committedConc)) + idsSeq := idsByPartitionValue(t, committedSeq) + idsConc := idsByPartitionValue(t, committedConc) + require.Len(t, idsConc, 8) for p := range 8 { - assert.Equal(t, int64(10), rowsByPartitionValue(t, committedConc)[fmt.Sprintf("p%d", p)]) + key := fmt.Sprintf("p%d", p) + assert.ElementsMatch(t, idsSeq[key], idsConc[key]) + assert.Len(t, idsConc[key], 10) } paths := manifestLiveDataPaths(t, committedConc) @@ -1706,11 +1713,13 @@ func TestRewriteDataFiles_MaxConcurrentGroupsGroupFailure(t *testing.T) { groupsAtomic := groupsByPartition(t, tblAtomic) require.Len(t, groupsAtomic, 4) fsAtomic.setFail(groupsAtomic[2].Tasks[0].File.FilePath(), injected) + beforeAtomicFiles := allParquetFiles(t, tblAtomic.Location()) txAtomic := tblAtomic.NewTransaction() _, err := txAtomic.RewriteDataFiles(t.Context(), groupsAtomic, table.RewriteDataFilesOptions{MaxConcurrentGroups: 4}) require.Error(t, err) assert.Contains(t, err.Error(), injected.Error()) + assert.ElementsMatch(t, beforeAtomicFiles, allParquetFiles(t, tblAtomic.Location())) fsPartial := &failOpenIO{} tblPartial := newMaxConcPartitionedTable(t, fsPartial) @@ -1780,42 +1789,44 @@ func (g *gateOpenIO) Open(name string) (iceio.File, error) { func TestRewriteDataFiles_MaxConcurrentGroupsContextCancel(t *testing.T) { for _, partial := range []bool{false, true} { - gate := &gateOpenIO{} - tbl := newMaxConcPartitionedTable(t, gate) - tbl = addMaxConcPartitions(t, tbl, 8, 1, 10) - groups := groupsByPartition(t, tbl) - require.Len(t, groups, 8) - gate.enable() - - ctx, cancel := context.WithCancel(t.Context()) - done := make(chan struct{}) - var rewriteErr error - go func() { - defer close(done) - tx := tbl.NewTransaction() - opts := table.RewriteDataFilesOptions{MaxConcurrentGroups: 4} - if partial { - opts.PartialProgress = true - opts.MaxCommits = 1 - } - _, rewriteErr = tx.RewriteDataFiles(ctx, groups, opts) - }() - - for range 4 { - select { - case <-gate.entered: - case <-done: - t.Fatalf("rewrite finished before 4 groups were in flight, err=%v", rewriteErr) - case <-t.Context().Done(): - t.Fatal("test context done while waiting for groups") + t.Run(fmt.Sprintf("partial=%v", partial), func(t *testing.T) { + gate := &gateOpenIO{} + tbl := newMaxConcPartitionedTable(t, gate) + tbl = addMaxConcPartitions(t, tbl, 8, 1, 10) + groups := groupsByPartition(t, tbl) + require.Len(t, groups, 8) + gate.enable() + + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan struct{}) + var rewriteErr error + go func() { + defer close(done) + tx := tbl.NewTransaction() + opts := table.RewriteDataFilesOptions{MaxConcurrentGroups: 4} + if partial { + opts.PartialProgress = true + opts.MaxCommits = 1 + } + _, rewriteErr = tx.RewriteDataFiles(ctx, groups, opts) + }() + + for range 4 { + select { + case <-gate.entered: + case <-done: + t.Fatalf("rewrite finished before 4 groups were in flight, err=%v", rewriteErr) + case <-t.Context().Done(): + t.Fatal("test context done while waiting for groups") + } } - } - cancel() - gate.releaseAll() - <-done - require.Error(t, rewriteErr) - assert.ErrorIs(t, rewriteErr, context.Canceled) - assert.Equal(t, ctx.Err(), rewriteErr) + cancel() + gate.releaseAll() + <-done + require.Error(t, rewriteErr) + assert.ErrorIs(t, rewriteErr, context.Canceled) + assert.Equal(t, ctx.Err(), rewriteErr) + }) } } @@ -1838,9 +1849,11 @@ func (f *countOpenFile) Close() error { type countOpenIO struct { iceio.LocalFS - mu sync.Mutex - cur int - peak int + mu sync.Mutex + cur int + peak int + barrier bool + overlap chan struct{} } func (c *countOpenIO) Open(name string) (iceio.File, error) { @@ -1854,8 +1867,18 @@ func (c *countOpenIO) Open(name string) (iceio.File, error) { if c.cur > c.peak { c.peak = c.cur } + if c.barrier && c.cur >= 2 { + select { + case <-c.overlap: + default: + close(c.overlap) + } + } + barrier, overlap := c.barrier, c.overlap c.mu.Unlock() - time.Sleep(20 * time.Millisecond) + if barrier { + <-overlap + } return &countOpenFile{File: f, owner: c}, nil } @@ -1868,6 +1891,16 @@ func (c *countOpenIO) reset() { defer c.mu.Unlock() c.cur = 0 c.peak = 0 + c.overlap = make(chan struct{}) +} + +func (c *countOpenIO) setBarrier(enabled bool) { + c.mu.Lock() + defer c.mu.Unlock() + c.barrier = enabled + if c.overlap == nil { + c.overlap = make(chan struct{}) + } } func (c *countOpenIO) getPeak() int { @@ -1879,28 +1912,31 @@ func (c *countOpenIO) getPeak() int { func TestRewriteDataFiles_MaxConcurrentGroupsLimitsInFlight(t *testing.T) { for _, maxConc := range []int{4, 0, 1} { - counter := &countOpenIO{} - tbl := newMaxConcPartitionedTable(t, counter) - tbl = addMaxConcPartitions(t, tbl, 8, 1, 10) - groups := groupsByPartition(t, tbl) - require.Len(t, groups, 8) - counter.reset() + t.Run(fmt.Sprintf("maxConc=%d", maxConc), func(t *testing.T) { + counter := &countOpenIO{} + tbl := newMaxConcPartitionedTable(t, counter) + tbl = addMaxConcPartitions(t, tbl, 8, 1, 10) + groups := groupsByPartition(t, tbl) + require.Len(t, groups, 8) + counter.reset() + counter.setBarrier(maxConc > 1) - tx := tbl.NewTransaction() - _, err := tx.RewriteDataFiles(t.Context(), groups, table.RewriteDataFilesOptions{ - MaxConcurrentGroups: maxConc, - GroupOptions: []table.CompactionGroupOption{table.WithCompactionScanConcurrency(1)}, - }) - require.NoError(t, err) - _, err = tx.Commit(t.Context()) - require.NoError(t, err) + tx := tbl.NewTransaction() + _, err := tx.RewriteDataFiles(t.Context(), groups, table.RewriteDataFilesOptions{ + MaxConcurrentGroups: maxConc, + GroupOptions: []table.CompactionGroupOption{table.WithCompactionScanConcurrency(1)}, + }) + require.NoError(t, err) + _, err = tx.Commit(t.Context()) + require.NoError(t, err) - peak := counter.getPeak() - if maxConc > 1 { - assert.LessOrEqual(t, peak, maxConc) - assert.GreaterOrEqual(t, peak, 2) - } else { - assert.Equal(t, 1, peak) - } + peak := counter.getPeak() + if maxConc > 1 { + assert.LessOrEqual(t, peak, maxConc) + assert.GreaterOrEqual(t, peak, 2) + } else { + assert.Equal(t, 1, peak) + } + }) } } From a283bb24d85e663d6226dffbcc41da5dc9be9bc9 Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Tue, 29 Sep 2026 21:13:55 -0400 Subject: [PATCH 08/10] test(codec): carry timestamp-ns partition round-trip coverage Carried along from the pre-rebase branch; the data_file_codec.go fix it originally accompanied merged separately as #2052. Keeps the Avro round-trip test and the ns entries in the partition-shape matrix. --- data_file_codec_test.go | 33 +++++++++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/data_file_codec_test.go b/data_file_codec_test.go index 7ca765849..3830af534 100644 --- a/data_file_codec_test.go +++ b/data_file_codec_test.go @@ -194,6 +194,37 @@ func TestMarshalAvroEntryDecimalPartitionRoundTrip(t *testing.T) { require.True(t, got.Equals(DecimalLiteral(want))) } +func TestMarshalAvroEntryTimestampPartitionRoundTrip(t *testing.T) { + // Micro cases run first. + // A nano type sharing their schema-cache key would decode as Timestamp instead of TimestampNano. + for _, tc := range []struct { + typ Type + want any + }{ + {TimestampType{}, Timestamp(1_700_000_000_000_000)}, + {TimestampTzType{}, Timestamp(1_700_000_000_000_000)}, + {TimestampNsType{}, TimestampNano(1_700_000_000_000_000_123)}, + {TimestampTzNsType{}, TimestampNano(1_700_000_000_000_000_123)}, + } { + t.Run(tc.typ.String(), func(t *testing.T) { + schema := NewSchema(0, NestedField{ID: 1, Name: "ts", Type: tc.typ}) + spec := NewPartitionSpecID(1, PartitionField{SourceIDs: []int{1}, FieldID: 1000, Name: "ts", Transform: IdentityTransform{}}) + + path := "s3://bucket/ns/tbl/data/ts.parquet" + builder, err := NewDataFileBuilder(spec, EntryContentData, path, ParquetFile, map[int]any{1000: tc.want}, nil, nil, 1, 1024) + require.NoError(t, err) + df, ok := builder.Build().(*dataFile) + require.True(t, ok) + + encoded, err := df.MarshalAvroEntry(spec, schema, 3) + require.NoError(t, err) + decoded, err := unmarshalAvroDataFileEntry(encoded, spec, schema, 3) + require.NoError(t, err) + require.Equal(t, tc.want, decoded.Partition()[1000]) + }) + } +} + // snapshotAvroFields returns a deep copy of every avro-tagged field on // d, keyed by field name. Slices, maps, byte arrays, and pointer // targets are reconstructed so the snapshot is fully independent of d @@ -357,6 +388,8 @@ func TestManifestEntrySchemaForMatchesPartitionAvroShape(t *testing.T) { TimeType{}, TimestampType{}, TimestampTzType{}, + TimestampNsType{}, + TimestampTzNsType{}, UUIDType{}, BooleanType{}, BinaryType{}, From 6c9a36ef7b29862d00980392fbf36d32c11e04d1 Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Tue, 29 Sep 2026 21:13:55 -0400 Subject: [PATCH 09/10] perf(table): harden atomic rewrite cleanup and concurrent error selection Wire cleanupAtomicRewriteOutputs into the rewrite.Commit failure branch so a staging failure after N parallel group writes removes every output, and open the table IO once up front so cancellation cannot fail the cleanup itself. cleanupAtomicRewriteOutputs now takes the open iceio.IO handle. Return the lowest-index failure not caused by sibling cancellation instead of errgroup's first-finished error, and log every group failure via slog. Cancellation now broadcasts over a plain WithCancel context rather than errgroup.WithContext so a sibling inherits context.Canceled instead of the failing group's error value through context.Cause. --- table/rewrite_data_files.go | 64 ++++++++++++++++++++++++++++--------- 1 file changed, 49 insertions(+), 15 deletions(-) diff --git a/table/rewrite_data_files.go b/table/rewrite_data_files.go index 174b02a82..cfa9da606 100644 --- a/table/rewrite_data_files.go +++ b/table/rewrite_data_files.go @@ -224,7 +224,10 @@ type RewriteDataFilesOptions struct { // Larger values run [ExecuteCompactionGroup] calls under a bounded // errgroup and apply their results in the original group order, so // manifests and [RewriteResult] are identical to a sequential run. - // The first error cancels the groups still running and is returned. + // A failure cancels the groups still running; the returned error + // is the lowest-index failure not caused by that cancellation, or + // the lowest-index failure when every group was canceled, so + // failures match a sequential run. Every group error is logged. // Peak record-pipeline memory is MaxConcurrentGroups times the // per-group bound stated on [WithCompactionArrowBatchSize]: // @@ -376,19 +379,25 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT rewrite := t.NewRewrite(opts.SnapshotProps) stagedDeleteFiles := make(map[string]struct{}) + fs, err := t.tbl.fsF(ctx) + if err != nil { + return result, fmt.Errorf("open table IO for atomic rewrite: %w", err) + } + + var applied []CompactionGroupResult if opts.MaxConcurrentGroups > 1 { results, err := executeCompactionGroups(ctx, t.tbl, groups, opts.GroupOptions, opts.MaxConcurrentGroups) if err != nil { - return result, cleanupAtomicRewriteOutputs(ctx, t.tbl, results, err) + return result, cleanupAtomicRewriteOutputs(fs, results, err) } + applied = results for _, gr := range results { applyAtomicGroupResult(rewrite, result, stagedDeleteFiles, gr) } } else { - var applied []CompactionGroupResult for _, group := range groups { if err := ctx.Err(); err != nil { - return result, cleanupAtomicRewriteOutputs(ctx, t.tbl, applied, err) + return result, cleanupAtomicRewriteOutputs(fs, applied, err) } if len(group.Tasks) == 0 { @@ -397,7 +406,7 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT gr, err := ExecuteCompactionGroup(ctx, t.tbl, group, opts.GroupOptions...) if err != nil { - return result, cleanupAtomicRewriteOutputs(ctx, t.tbl, append(applied, gr), err) + return result, cleanupAtomicRewriteOutputs(fs, append(applied, gr), err) } applied = append(applied, gr) @@ -430,7 +439,7 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT } if err := rewrite.Commit(ctx); err != nil { - return result, fmt.Errorf("commit compaction: %w", err) + return result, cleanupAtomicRewriteOutputs(fs, applied, fmt.Errorf("commit compaction: %w", err)) } return result, nil @@ -450,11 +459,7 @@ func applyAtomicGroupResult(rewrite *RewriteFiles, result *RewriteResult, staged } } -func cleanupAtomicRewriteOutputs(ctx context.Context, tbl *Table, results []CompactionGroupResult, cause error) error { - fs, err := tbl.fsF(ctx) - if err != nil { - return errors.Join(cause, fmt.Errorf("open table IO to clean up atomic rewrite outputs: %w", err)) - } +func cleanupAtomicRewriteOutputs(fs iceio.IO, results []CompactionGroupResult, cause error) error { if err := cleanupCompactionOutputs(fs, results); err != nil { return errors.Join(cause, fmt.Errorf("clean up atomic rewrite outputs: %w", err)) } @@ -470,26 +475,55 @@ func executeCompactionGroups(ctx context.Context, tbl *Table, groups []Compactio if limit < 1 { limit = 1 } - g, gctx := errgroup.WithContext(ctx) + var g errgroup.Group g.SetLimit(limit) + runCtx, cancelRuns := context.WithCancel(ctx) + defer cancelRuns() results := make([]CompactionGroupResult, len(groups)) + groupErrs := make([]error, len(groups)) for i, group := range groups { if len(group.Tasks) == 0 { continue } g.Go(func() error { - gr, err := ExecuteCompactionGroup(gctx, tbl, group, groupOpts...) + gr, err := ExecuteCompactionGroup(runCtx, tbl, group, groupOpts...) results[i] = gr + groupErrs[i] = err + if err != nil { + slog.Warn("compaction group failed", "index", i, "err", err) + cancelRuns() + } return err }) } if err := g.Wait(); err != nil { - if ctxErr := ctx.Err(); ctxErr != nil && (errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded)) { + var firstErr, firstNonContextErr error + for _, groupErr := range groupErrs { + if groupErr == nil { + continue + } + if firstErr == nil { + firstErr = groupErr + } + if !errors.Is(groupErr, context.Canceled) && !errors.Is(groupErr, context.DeadlineExceeded) { + firstNonContextErr = groupErr + + break + } + } + selected := firstNonContextErr + if selected == nil { + selected = firstErr + } + if selected == nil { + selected = err + } + if ctxErr := ctx.Err(); ctxErr != nil && (errors.Is(selected, context.Canceled) || errors.Is(selected, context.DeadlineExceeded)) { return results, ctxErr } - return results, err + return results, selected } return results, nil From 11161876920da47543ad650666948d74eae1e87b Mon Sep 17 00:00:00 2001 From: KranzL <50032317+KranzL@users.noreply.github.com> Date: Tue, 29 Sep 2026 21:13:55 -0400 Subject: [PATCH 10/10] test(table): pin atomic cleanup and concurrent error selection CommitFailureCleansOutputs fails rewrite.Commit through an unsupported content type and asserts the before/after parquet set is unchanged for both the concurrent and sequential paths. CleanupAfterCancelUsesOpenFS runs the table IO factory that fails on a done context and asserts full cleanup after mid-run cancellation on both paths. FailureReturnsLowest- IndexError orders two injected failures so the higher index finishes first and asserts the lower index wins while both are logged. --- table/rewrite_data_files_test.go | 234 ++++++++++++++++++++++++++++++- 1 file changed, 233 insertions(+), 1 deletion(-) diff --git a/table/rewrite_data_files_test.go b/table/rewrite_data_files_test.go index a5b747979..eae48b51a 100644 --- a/table/rewrite_data_files_test.go +++ b/table/rewrite_data_files_test.go @@ -18,16 +18,19 @@ package table_test import ( + "bytes" "context" "errors" "fmt" "io/fs" + "log/slog" "os" "path/filepath" "slices" "strings" "sync" "testing" + "time" "github.com/apache/arrow-go/v18/arrow" "github.com/apache/arrow-go/v18/arrow/array" @@ -1427,6 +1430,12 @@ func appendEqualityDelete(t *testing.T, tbl *table.Table, equalityFieldIDs []int func newMaxConcPartitionedTable(t *testing.T, fs iceio.IO) *table.Table { t.Helper() + return newMaxConcPartitionedTableWithFSF(t, func(context.Context) (iceio.IO, error) { return fs, nil }) +} + +func newMaxConcPartitionedTableWithFSF(t *testing.T, fsF func(context.Context) (iceio.IO, error)) *table.Table { + t.Helper() + location := filepath.ToSlash(t.TempDir()) schema := iceberg.NewSchema(0, iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64, Required: true}, @@ -1444,7 +1453,7 @@ func newMaxConcPartitionedTable(t *testing.T, fs iceio.IO) *table.Table { return table.New( table.Identifier{"db", "max_conc_test"}, meta, location+"/metadata/v1.metadata.json", - func(context.Context) (iceio.IO, error) { return fs, nil }, + fsF, cat, ) } @@ -1742,6 +1751,229 @@ func TestRewriteDataFiles_MaxConcurrentGroupsGroupFailure(t *testing.T) { assert.ElementsMatch(t, beforeFiles, allParquetFiles(t, tblPartial.Location())) } +type badContentFile struct { + iceberg.DataFile +} + +func (badContentFile) ContentType() iceberg.ManifestEntryContent { + return iceberg.ManifestEntryContent(99) +} + +func TestRewriteDataFiles_MaxConcurrentGroupsCommitFailureCleansOutputs(t *testing.T) { + for _, maxConc := range []int{4, 0} { + t.Run(fmt.Sprintf("maxConc=%d", maxConc), func(t *testing.T) { + tbl := newMaxConcPartitionedTable(t, iceio.LocalFS{}) + tbl = addMaxConcPartitions(t, tbl, 4, 1, 5) + groups := groupsByPartition(t, tbl) + require.Len(t, groups, 4) + before := allParquetFiles(t, tbl.Location()) + + tx := tbl.NewTransaction() + _, err := tx.RewriteDataFiles(t.Context(), groups, table.RewriteDataFilesOptions{ + MaxConcurrentGroups: maxConc, + ExtraDeleteFilesToRemove: []iceberg.DataFile{badContentFile{groups[0].Tasks[0].File}}, + }) + require.Error(t, err) + assert.Contains(t, err.Error(), "unsupported content type") + assert.ElementsMatch(t, before, allParquetFiles(t, tbl.Location())) + }) + } +} + +type blockPathIO struct { + iceio.LocalFS + mu sync.Mutex + substr string + entered chan struct{} + release chan struct{} +} + +func (b *blockPathIO) Open(name string) (iceio.File, error) { + b.mu.Lock() + substr, entered, release := b.substr, b.entered, b.release + b.mu.Unlock() + if substr != "" && strings.Contains(name, substr) { + select { + case entered <- struct{}{}: + default: + } + <-release + } + + return b.LocalFS.Open(name) +} + +func cancelOnDoneFSF(t *testing.T, fs iceio.IO) func(context.Context) (iceio.IO, error) { + t.Helper() + + return func(ctx context.Context) (iceio.IO, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + + return fs, nil + } +} + +func TestRewriteDataFiles_MaxConcurrentGroupsCleanupAfterCancelUsesOpenFS(t *testing.T) { + t.Run("sequential", func(t *testing.T) { + blocker := &blockPathIO{entered: make(chan struct{}, 1), release: make(chan struct{})} + tbl := newMaxConcPartitionedTableWithFSF(t, cancelOnDoneFSF(t, blocker)) + tbl = addMaxConcPartitions(t, tbl, 4, 1, 5) + groups := groupsByPartition(t, tbl) + require.Len(t, groups, 4) + blocker.mu.Lock() + blocker.substr = groups[1].Tasks[0].File.FilePath() + blocker.mu.Unlock() + before := allParquetFiles(t, tbl.Location()) + + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan struct{}) + var rewriteErr error + go func() { + defer close(done) + tx := tbl.NewTransaction() + _, rewriteErr = tx.RewriteDataFiles(ctx, groups, table.RewriteDataFilesOptions{}) + }() + + select { + case <-blocker.entered: + case <-done: + t.Fatalf("rewrite finished before group 1 started, err=%v", rewriteErr) + case <-t.Context().Done(): + t.Fatal("test context done while waiting for group 1") + } + cancel() + close(blocker.release) + <-done + require.Error(t, rewriteErr) + assert.ErrorIs(t, rewriteErr, context.Canceled) + assert.ElementsMatch(t, before, allParquetFiles(t, tbl.Location())) + }) + + t.Run("concurrent", func(t *testing.T) { + blocker := &blockPathIO{entered: make(chan struct{}, 1), release: make(chan struct{})} + tbl := newMaxConcPartitionedTableWithFSF(t, cancelOnDoneFSF(t, blocker)) + tbl = addMaxConcPartitions(t, tbl, 8, 1, 10) + groups := groupsByPartition(t, tbl) + require.Len(t, groups, 8) + blocker.mu.Lock() + blocker.substr = groups[7].Tasks[0].File.FilePath() + blocker.mu.Unlock() + before := allParquetFiles(t, tbl.Location()) + + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan struct{}) + var rewriteErr error + go func() { + defer close(done) + tx := tbl.NewTransaction() + _, rewriteErr = tx.RewriteDataFiles(ctx, groups, table.RewriteDataFilesOptions{MaxConcurrentGroups: 4}) + }() + + deadline := time.Now().Add(time.Minute) + for len(allParquetFiles(t, tbl.Location())) == len(before) { + select { + case <-done: + t.Fatalf("rewrite finished before any group wrote output, err=%v", rewriteErr) + default: + } + if time.Now().After(deadline) { + cancel() + close(blocker.release) + <-done + t.Fatal("timed out waiting for a group to write output") + } + time.Sleep(5 * time.Millisecond) + } + cancel() + close(blocker.release) + <-done + require.Error(t, rewriteErr) + assert.ErrorIs(t, rewriteErr, context.Canceled) + assert.ElementsMatch(t, before, allParquetFiles(t, tbl.Location())) + }) +} + +type orderFailIO struct { + iceio.LocalFS + mu sync.Mutex + slowSubstr string + slowErr error + fastSubstr string + fastErr error + slowEntered chan struct{} + fastFailed chan struct{} + slowSignaled bool + fastSignaled bool +} + +func (o *orderFailIO) Open(name string) (iceio.File, error) { + o.mu.Lock() + slow := o.slowSubstr != "" && strings.Contains(name, o.slowSubstr) + fast := o.fastSubstr != "" && strings.Contains(name, o.fastSubstr) + if slow && !o.slowSignaled { + o.slowSignaled = true + close(o.slowEntered) + } + slowErr, fastErr := o.slowErr, o.fastErr + slowEntered, fastFailed := o.slowEntered, o.fastFailed + o.mu.Unlock() + switch { + case slow: + <-fastFailed + + return nil, slowErr + case fast: + <-slowEntered + o.mu.Lock() + if !o.fastSignaled { + o.fastSignaled = true + close(fastFailed) + } + o.mu.Unlock() + + return nil, fastErr + default: + return o.LocalFS.Open(name) + } +} + +func TestRewriteDataFiles_MaxConcurrentGroupsFailureReturnsLowestIndexError(t *testing.T) { + slowErr := errors.New("injected slow group failure") + fastErr := errors.New("injected fast group failure") + + ordered := &orderFailIO{slowEntered: make(chan struct{}), fastFailed: make(chan struct{})} + tbl := newMaxConcPartitionedTable(t, ordered) + tbl = addMaxConcPartitions(t, tbl, 4, 1, 5) + groups := groupsByPartition(t, tbl) + require.Len(t, groups, 4) + ordered.mu.Lock() + ordered.slowSubstr = groups[1].Tasks[0].File.FilePath() + ordered.slowErr = slowErr + ordered.fastSubstr = groups[3].Tasks[0].File.FilePath() + ordered.fastErr = fastErr + ordered.mu.Unlock() + before := allParquetFiles(t, tbl.Location()) + + var buf bytes.Buffer + origLogger := slog.Default() + slog.SetDefault(slog.New(slog.NewTextHandler(&buf, &slog.HandlerOptions{Level: slog.LevelWarn}))) + defer slog.SetDefault(origLogger) + + tx := tbl.NewTransaction() + _, err := tx.RewriteDataFiles(t.Context(), groups, table.RewriteDataFilesOptions{MaxConcurrentGroups: 4}) + require.Error(t, err) + assert.Contains(t, err.Error(), slowErr.Error()) + assert.NotContains(t, err.Error(), fastErr.Error()) + assert.ElementsMatch(t, before, allParquetFiles(t, tbl.Location())) + + logged := buf.String() + assert.Contains(t, logged, "compaction group failed") + assert.Contains(t, logged, slowErr.Error()) + assert.Contains(t, logged, fastErr.Error()) +} + type gateOpenIO struct { iceio.LocalFS mu sync.Mutex