From 4ef3988e724975cbfa0faf9ef653fff45881b156 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Sat, 29 Aug 2026 15:21:10 +0200 Subject: [PATCH 1/3] perf(table): parallelize non-bulk purge deletion Signed-off-by: Minh Vu --- table/orphan_cleanup.go | 152 +++++++++++++++------- table/orphan_cleanup_bench_test.go | 44 +++++++ table/orphan_cleanup_test.go | 194 +++++++++++++++++++++++++++++ 3 files changed, 342 insertions(+), 48 deletions(-) diff --git a/table/orphan_cleanup.go b/table/orphan_cleanup.go index b20b8b033..294c7cda1 100644 --- a/table/orphan_cleanup.go +++ b/table/orphan_cleanup.go @@ -688,14 +688,27 @@ func deleteFiles(ctx context.Context, fs iceio.IO, orphanFiles []string, cfg *or } } - if cfg.maxConcurrency == 1 { - return deleteFilesSequential(fs, orphanFiles, cfg) + if cfg.maxConcurrency <= 1 { + return deleteFilesSequential(ctx, fs, orphanFiles, cfg) } - return deleteFilesParallel(fs, orphanFiles, cfg) + deleteFunc := fs.Remove + if cfg.deleteFunc != nil { + deleteFunc = cfg.deleteFunc + } + + return deleteFilesParallel( + ctx, + orphanFiles, + cfg.maxConcurrency, + deleteFunc, + func(file string, err error) error { + return fmt.Errorf("failed to delete orphan file %s: %w", file, err) + }, + ) } -func deleteFilesSequential(fs iceio.IO, orphanFiles []string, cfg *orphanCleanupConfig) ([]string, error) { +func deleteFilesSequential(ctx context.Context, fs iceio.IO, orphanFiles []string, cfg *orphanCleanupConfig) ([]string, error) { var deletedFiles []string deleteFunc := fs.Remove @@ -705,6 +718,11 @@ func deleteFilesSequential(fs iceio.IO, orphanFiles []string, cfg *orphanCleanup var result error for _, file := range orphanFiles { + if err := ctx.Err(); err != nil { + result = errors.Join(result, err) + + break + } if err := deleteFunc(file); err != nil { result = errors.Join(result, fmt.Errorf("failed to delete orphan file %s: %w", file, err)) @@ -716,55 +734,85 @@ func deleteFilesSequential(fs iceio.IO, orphanFiles []string, cfg *orphanCleanup return deletedFiles, result } -func deleteFilesParallel(fs iceio.IO, orphanFiles []string, cfg *orphanCleanupConfig) ([]string, error) { - deleteFunc := fs.Remove - if cfg.deleteFunc != nil { - deleteFunc = cfg.deleteFunc +func deleteFilesParallel( + ctx context.Context, + files []string, + maxConcurrency int, + deleteFunc func(string) error, + wrapError func(string, error) error, +) ([]string, error) { + workers := min(max(maxConcurrency, 1), len(files)) + jobs := make(chan int) + deleted := make([]bool, len(files)) + deleteErrors := make([]error, len(files)) + + var cancellationErr error + var cancellationOnce sync.Once + recordCancellation := func(err error) { + cancellationOnce.Do(func() { + cancellationErr = err + }) } - in := make(chan string, cfg.maxConcurrency) - out := make(chan string, cfg.maxConcurrency) - errList := make([][]error, cfg.maxConcurrency) - - go func() { - defer close(in) - for _, file := range orphanFiles { - in <- file - } - }() - var wg sync.WaitGroup - wg.Add(cfg.maxConcurrency) - for i := range cfg.maxConcurrency { - go func(workerID int) { + wg.Add(workers) + for range workers { + go func() { defer wg.Done() - for file := range in { - if err := deleteFunc(file); err != nil { - errList[workerID] = append(errList[workerID], fmt.Errorf("failed to delete orphan file %s: %w", file, err)) - } else { - out <- file + for { + select { + case <-ctx.Done(): + recordCancellation(ctx.Err()) + + return + case index, ok := <-jobs: + if !ok { + return + } + if err := ctx.Err(); err != nil { + recordCancellation(err) + + return + } + + if err := deleteFunc(files[index]); err != nil { + deleteErrors[index] = wrapError(files[index], err) + } else { + deleted[index] = true + } } } - }(i) + }() } - go func() { - wg.Wait() - close(out) - }() +send: + for index := range files { + select { + case jobs <- index: + case <-ctx.Done(): + recordCancellation(ctx.Err()) - deletedFiles := make([]string, 0, len(orphanFiles)) - for file := range out { - deletedFiles = append(deletedFiles, file) + break send + } } + close(jobs) + wg.Wait() - var allErrors []error - for _, workerErrors := range errList { - allErrors = append(allErrors, workerErrors...) + deletedFiles := make([]string, 0, len(files)) + allErrors := make([]error, 0) + for index, file := range files { + if deleted[index] { + deletedFiles = append(deletedFiles, file) + } + if deleteErrors[index] != nil { + allErrors = append(allErrors, deleteErrors[index]) + } + } + if cancellationErr != nil { + allErrors = append(allErrors, cancellationErr) } - err := errors.Join(allErrors...) - return deletedFiles, err + return deletedFiles, errors.Join(allErrors...) } // normalizeFilePath normalizes file paths for comparison by handling different @@ -1310,15 +1358,23 @@ func (t Table) PurgeFiles(ctx context.Context) error { errs = append(errs, fmt.Errorf("bulk deletion failed: %w", bulkErr)) } } else { - for _, file := range files { - if err := ctx.Err(); err != nil { - errs = append(errs, err) + _, removeErr := deleteFilesParallel( + ctx, + files, + runtime.GOMAXPROCS(0), + func(file string) error { + if err := fs.Remove(file); err != nil && !os.IsNotExist(err) { + return err + } - break - } - if rmErr := fs.Remove(file); rmErr != nil && !os.IsNotExist(rmErr) { - errs = append(errs, fmt.Errorf("failed to remove %s: %w", file, rmErr)) - } + return nil + }, + func(file string, err error) error { + return fmt.Errorf("failed to remove %s: %w", file, err) + }, + ) + if removeErr != nil { + errs = append(errs, removeErr) } } } diff --git a/table/orphan_cleanup_bench_test.go b/table/orphan_cleanup_bench_test.go index 52b1cba13..70ff425db 100644 --- a/table/orphan_cleanup_bench_test.go +++ b/table/orphan_cleanup_bench_test.go @@ -18,8 +18,10 @@ package table import ( + "context" "fmt" "testing" + "time" ) var orphanCleanupBenchmarkSink string @@ -63,3 +65,45 @@ func BenchmarkApplyURIEquivalence(b *testing.B) { }) } } + +// BenchmarkPurgeFilesNonBulkDeletion measures the bounded fallback used when +// the filesystem does not implement BulkRemovableIO. The delay models the +// round trip to a remote object store; the zero-delay cases show the worker +// pool overhead on local-style deletes. +func BenchmarkPurgeFilesNonBulkDeletion(b *testing.B) { + for _, fileCount := range []int{100, 1_000, 10_000} { + files := make([]string, fileCount) + for i := range files { + files[i] = fmt.Sprintf("s3://bucket/table/data/file-%d.parquet", i) + } + + for _, delay := range []time.Duration{0, 100 * time.Microsecond, time.Millisecond} { + for _, concurrency := range []int{1, 4, 16} { + b.Run(fmt.Sprintf("files=%d/delay=%s/concurrency=%d", fileCount, delay, concurrency), func(b *testing.B) { + deleteFunc := func(string) error { + time.Sleep(delay) + + return nil + } + + b.ReportAllocs() + b.ReportMetric(float64(fileCount), "files/op") + b.ResetTimer() + for b.Loop() { + deleted, err := deleteFilesParallel( + context.Background(), + files, + concurrency, + deleteFunc, + func(string, error) error { return nil }, + ) + if err != nil { + b.Fatal(err) + } + orphanCleanupBenchmarkSink = deleted[len(deleted)-1] + } + }) + } + } + } +} diff --git a/table/orphan_cleanup_test.go b/table/orphan_cleanup_test.go index 7592261db..a1243e2d7 100644 --- a/table/orphan_cleanup_test.go +++ b/table/orphan_cleanup_test.go @@ -22,9 +22,11 @@ import ( "errors" "fmt" stdfs "io/fs" + "os" "path/filepath" "runtime" "strings" + "sync" "testing" "time" @@ -1567,6 +1569,7 @@ func TestDeleteFilesEmpty(t *testing.T) { // mockPlainIO implements only IO, without optional delete or listing capabilities. type mockPlainIO struct { + mu sync.Mutex removed []string } @@ -1575,6 +1578,9 @@ func (m *mockPlainIO) Open(string) (io.File, error) { } func (m *mockPlainIO) Remove(name string) error { + m.mu.Lock() + defer m.mu.Unlock() + m.removed = append(m.removed, name) return nil @@ -1616,6 +1622,194 @@ func (m mockFileInfo) ModTime() time.Time { return m.modTime } func (m mockFileInfo) IsDir() bool { return m.mode.IsDir() } func (m mockFileInfo) Sys() any { return nil } +type purgeDeleteTrackingIO struct { + mockListableIO + + targetActive int + reached chan struct{} + release chan struct{} + + mu sync.Mutex + active int + maxActive int + removed []string + reachedOnce sync.Once + releaseOnce sync.Once +} + +func (m *purgeDeleteTrackingIO) Remove(name string) error { + m.mu.Lock() + m.active++ + if m.active > m.maxActive { + m.maxActive = m.active + } + if m.targetActive > 0 && m.active >= m.targetActive { + m.reachedOnce.Do(func() { close(m.reached) }) + } + m.mu.Unlock() + + if m.targetActive > 0 { + <-m.release + } + + m.mu.Lock() + m.active-- + m.removed = append(m.removed, name) + m.mu.Unlock() + + return nil +} + +func (m *purgeDeleteTrackingIO) MaxActive() int { + m.mu.Lock() + defer m.mu.Unlock() + + return m.maxActive +} + +func (m *purgeDeleteTrackingIO) Release() { + m.releaseOnce.Do(func() { close(m.release) }) +} + +func TestDeleteFilesParallelCollectsPurgeErrors(t *testing.T) { + const ( + firstPath = "s3://bucket/table/first.parquet" + missingPath = "s3://bucket/table/missing.parquet" + slowPath = "s3://bucket/table/slow.parquet" + fastPath = "s3://bucket/table/fast.parquet" + ) + + slowErr := errors.New("slow removal failed") + fastErr := errors.New("fast removal failed") + files := []string{firstPath, missingPath, slowPath, fastPath} + var mu sync.Mutex + calls := make(map[string]int) + + deleted, err := deleteFilesParallel( + context.Background(), + files, + 4, + func(path string) error { + mu.Lock() + calls[path]++ + mu.Unlock() + + var err error + switch path { + case missingPath: + err = stdfs.ErrNotExist + case slowPath: + time.Sleep(10 * time.Millisecond) + + err = slowErr + case fastPath: + err = fastErr + } + if os.IsNotExist(err) { + return nil + } + + return err + }, + func(path string, err error) error { + return fmt.Errorf("failed to remove %s: %w", path, err) + }, + ) + + assert.Equal(t, []string{firstPath, missingPath}, deleted) + require.ErrorIs(t, err, slowErr) + require.ErrorIs(t, err, fastErr) + assert.NotContains(t, err.Error(), missingPath) + assert.Less(t, strings.Index(err.Error(), slowPath), strings.Index(err.Error(), fastPath)) + + mu.Lock() + assert.Equal(t, map[string]int{ + firstPath: 1, + missingPath: 1, + slowPath: 1, + fastPath: 1, + }, calls) + mu.Unlock() +} + +func TestDeleteFilesParallelStopsQueuedWorkOnCancellation(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + called := false + _, err := deleteFilesParallel( + ctx, + []string{"s3://bucket/table/file.parquet"}, + 4, + func(string) error { + called = true + + return nil + }, + func(path string, err error) error { + return fmt.Errorf("failed to remove %s: %w", path, err) + }, + ) + + require.ErrorIs(t, err, context.Canceled) + assert.False(t, called) +} + +func TestPurgeFilesDeletesNonBulkFilesConcurrently(t *testing.T) { + maxWorkers := runtime.GOMAXPROCS(0) + if maxWorkers < 2 { + t.Skip("parallel deletion requires at least two workers") + } + + const fileCount = 16 + entries := make([]mockWalkEntry, 0, fileCount) + for i := range fileCount { + entries = append(entries, mockWalkEntry{ + path: fmt.Sprintf("s3://bucket/table/data/file-%02d.parquet", i), + info: mockFileInfo{name: fmt.Sprintf("file-%02d.parquet", i)}, + }) + } + + fsys := &purgeDeleteTrackingIO{ + mockListableIO: mockListableIO{entries: entries}, + targetActive: 2, + reached: make(chan struct{}), + release: make(chan struct{}), + } + + meta, err := NewMetadata( + iceberg.NewSchema(0), + iceberg.UnpartitionedSpec, + UnsortedSortOrder, + "s3://bucket/table", + iceberg.Properties{}, + ) + require.NoError(t, err) + tbl := New( + Identifier{"db", "tbl"}, + meta, + "s3://bucket/table/metadata/v1.metadata.json", + testFSF(fsys), + nil, + ) + + done := make(chan error, 1) + go func() { done <- tbl.PurgeFiles(context.Background()) }() + + select { + case <-fsys.reached: + case <-time.After(5 * time.Second): + fsys.Release() + <-done + t.Fatal("non-bulk purge deletion did not reach two concurrent removals") + } + fsys.Release() + + require.NoError(t, <-done) + assert.Greater(t, fsys.MaxActive(), 1) + assert.LessOrEqual(t, fsys.MaxActive(), maxWorkers) +} + func TestPurgeFilesSkipsDataFilesForMalformedGCEnabled(t *testing.T) { const orphanDataPath = "s3://bucket/table/data/orphan.parquet" From ddd4db916e6364dc40f0427732fc4348b2eb5a3f Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Sun, 30 Aug 2026 23:18:27 +0200 Subject: [PATCH 2/3] test(catalog): synchronize parallel purge callbacks Signed-off-by: Minh Vu --- catalog/glue/glue_test.go | 23 +++++++++++++---------- catalog/hive/hive_test.go | 23 +++++++++++++---------- 2 files changed, 26 insertions(+), 20 deletions(-) diff --git a/catalog/glue/glue_test.go b/catalog/glue/glue_test.go index 46d194f15..92c45d240 100644 --- a/catalog/glue/glue_test.go +++ b/catalog/glue/glue_test.go @@ -28,6 +28,7 @@ import ( "path/filepath" "strconv" "strings" + "sync/atomic" "testing" "time" @@ -864,16 +865,18 @@ func TestGluePurgeTableSwallowsPurgeFilesError(t *testing.T) { assert := require.New(t) ctx := context.Background() const scheme = "gluepurgefail" - dropCalled := false - removeBeforeDrop := false - removeCalls := 0 + var ( + dropCalled atomic.Bool + removeBeforeDrop atomic.Bool + removeCalls atomic.Int64 + ) failingFS := failRemoveIO{ MemFS: iceio.NewMemFS(), err: errGluePurgeRemove, onRemove: func() { - removeCalls++ - if !dropCalled { - removeBeforeDrop = true + removeCalls.Add(1) + if !dropCalled.Load() { + removeBeforeDrop.Store(true) } }, } @@ -899,15 +902,15 @@ func TestGluePurgeTableSwallowsPurgeFilesError(t *testing.T) { DatabaseName: aws.String("test_database"), Name: aws.String("test_table"), }, mock.Anything).Run(func(mock.Arguments) { - dropCalled = true + dropCalled.Store(true) }).Return(&glue.DeleteTableOutput{}, nil).Once() glueCatalog := &Catalog{glueSvc: mockGlueSvc} assert.NoError(glueCatalog.PurgeTable(ctx, TableIdentifier("test_database", "test_table"))) - assert.True(dropCalled) - assert.Positive(removeCalls) - assert.False(removeBeforeDrop, "PurgeTable should drop the catalog entry before removing files") + assert.True(dropCalled.Load()) + assert.Positive(removeCalls.Load()) + assert.False(removeBeforeDrop.Load(), "PurgeTable should drop the catalog entry before removing files") file, err := failingFS.Open(dataFile) assert.NoError(err, "data file should remain when FileIO remove fails") assert.NotNil(file) diff --git a/catalog/hive/hive_test.go b/catalog/hive/hive_test.go index f4f3095de..85ab555e7 100644 --- a/catalog/hive/hive_test.go +++ b/catalog/hive/hive_test.go @@ -25,6 +25,7 @@ import ( "os" "path/filepath" "strings" + "sync/atomic" "testing" "github.com/apache/iceberg-go" @@ -916,16 +917,18 @@ func TestHivePurgeTableSwallowsPurgeFilesError(t *testing.T) { assert := require.New(t) ctx := context.Background() const scheme = "hivepurgefail" - dropCalled := false - removeBeforeDrop := false - removeCalls := 0 + var ( + dropCalled atomic.Bool + removeBeforeDrop atomic.Bool + removeCalls atomic.Int64 + ) failingFS := failRemoveIO{ MemFS: iceio.NewMemFS(), err: errHivePurgeRemove, onRemove: func() { - removeCalls++ - if !dropCalled { - removeBeforeDrop = true + removeCalls.Add(1) + if !dropCalled.Load() { + removeBeforeDrop.Store(true) } }, } @@ -948,15 +951,15 @@ func TestHivePurgeTableSwallowsPurgeFilesError(t *testing.T) { Return(hiveTable, nil).Twice() mockClient.On("DropTable", mock.Anything, "test_database", "test_table", false). Run(func(mock.Arguments) { - dropCalled = true + dropCalled.Store(true) }).Return(nil).Once() hiveCatalog := NewCatalogWithClient(mockClient, iceberg.Properties{}) assert.NoError(hiveCatalog.PurgeTable(ctx, TableIdentifier("test_database", "test_table"))) - assert.True(dropCalled) - assert.Positive(removeCalls) - assert.False(removeBeforeDrop, "PurgeTable should drop the catalog entry before removing files") + assert.True(dropCalled.Load()) + assert.Positive(removeCalls.Load()) + assert.False(removeBeforeDrop.Load(), "PurgeTable should drop the catalog entry before removing files") file, err := failingFS.Open(dataFile) assert.NoError(err, "data file should remain when FileIO remove fails") assert.NotNil(file) From 1cd87686047ad1df85de5078449b63a3bbb60624 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Sun, 30 Aug 2026 23:28:52 +0200 Subject: [PATCH 3/3] 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) + } +}