This is an automated email from the ASF dual-hosted git repository.
laskoviymishka pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-go.git
The following commit(s) were added to refs/heads/main by this push:
new 9cd5a0559 refactor(table)!: extract orphan cleanup into
table/maintenance (#2086)
9cd5a0559 is described below
commit 9cd5a0559ca0ca5e947fcf4da94f52972c7774b4
Author: Andrei Tserakhau <[email protected]>
AuthorDate: Mon Oct 5 23:02:54 2026 +0200
refactor(table)!: extract orphan cleanup into table/maintenance (#2086)
* refactor(table)!: extract orphan cleanup into table/maintenance
This is Slice 1 of issue #2083, breaking pre-v1 relocation of the orphan
cleanup
service from the table package into a dedicated table/maintenance package.
Changes:
- Create new table/maintenance package with Service facade
- Move DeleteOrphanFiles, PlanOrphanFiles, ExecuteOrphanCleanup, and
PurgeFiles
methods from table.Table to maintenance.Service
- Move exported types: OrphanCleanupOption, OrphanCleanupResult,
OrphanCleanupPlan,
OrphanFile, PrefixMismatchMode (and constants)
- Move helper functions: WithLocation, WithFilesOlderThan, WithDryRun,
WithDeleteFunc, WithCleanupMaxConcurrency, WithPrefixMismatchMode,
WithEqualSchemes, WithEqualAuthorities
- Export IsGCEnabled from table package (was unexported isGCEnabled)
- Update all callers in catalog/ and cmd/iceberg/ to use new API:
maintenance.New(tbl).DeleteOrphanFiles(...) instead of
tbl.DeleteOrphanFiles(...)
- Move tests to table/maintenance package
- Update IsGCEnabled references in transaction.go and updates.go
The new public entry shape maintains ergonomic parity with the original API.
To use orphan cleanup: create a maintenance service with
maintenance.New(table),
then call its methods.
BREAKING: table.DeleteOrphanFiles, table.PlanOrphanFiles,
table.ExecuteOrphanCleanup,
and table.PurgeFiles are no longer available on table.Table. Use
table/maintenance
package instead.
* style(table): gofmt the moved orphan-cleanup test
table/maintenance/orphan_cleanup_test.go was not gofmt-clean after the move
(a
block comment needed re-indentation), which failed the golangci-lint gofmt
check
in CI. Formatting only, no behavior change.
* refactor(table/maintenance): use free functions and restore dropped tests
Address review feedback on the orphan-cleanup extraction:
- Replace the Service facade with free functions that take the table as
their
first argument (maintenance.DeleteOrphanFiles(ctx, tbl, ...)), matching
table/compaction, and update all callers in catalog/ and cmd/iceberg/.
This
avoids freezing a stateful constructor at v1; a deprecation shim on
table.Table is impossible anyway because table/maintenance imports table.
- Rename the four generic options to the WithCleanup* prefix
(WithCleanupLocation/DryRun/DeleteFunc/FilesOlderThan) for consistency
with
WithCleanupMaxConcurrency before the names freeze.
- Restore the tests and benchmarks dropped in the move: the statistics- and
tombstone-reference tests, the manifest-dedup and concurrency-limit tests,
the negative-/zero-age retention checks, and the PurgeFiles and
manifest-list
benchmarks, porting their shared IO and manifest helpers into the package.
- Document table.IsGCEnabled's parsing rules now that it is exported, and
add a
package doc comment.
- Fix website/src/api.md and releases.md (with a moved-symbols compatibility
note), and add a `make vet` target running `go vet -tags=integration
./...`
(wired into CI) so the moved //go:build integration test stays compiled.
* docs(table/maintenance): document OrphanCleanupOption and the
WithCleanup* options
Add one-line doc comments with defaults to OrphanCleanupOption,
WithCleanupLocation (table data location), WithCleanupFilesOlderThan (72h),
and WithCleanupDryRun (false), now that the WithCleanup* names are fixed
for v1.
---
.github/workflows/go-ci.yml | 3 +
Makefile | 9 +-
catalog/glue/glue.go | 3 +-
catalog/hadoop/hadoop.go | 3 +-
catalog/hive/hive.go | 3 +-
catalog/sql/sql.go | 3 +-
cmd/iceberg/clean_orphan_files.go | 15 +-
cmd/iceberg/clean_orphan_files_test.go | 13 +-
table/maintenance/manifest_test_helpers_test.go | 243 ++++++++++++++++
table/{ => maintenance}/orphan_cleanup.go | 93 ++++---
.../{ => maintenance}/orphan_cleanup_bench_test.go | 25 +-
.../orphan_cleanup_integration_test.go | 43 +--
table/{ => maintenance}/orphan_cleanup_test.go | 304 +++++++++++----------
table/maintenance/retention_validation_test.go | 38 +++
table/properties.go | 10 +-
table/retention_validation_test.go | 10 -
table/transaction.go | 2 +-
table/transaction_internal_test.go | 2 +-
table/updates.go | 2 +-
website/src/api.md | 12 +-
website/src/releases.md | 9 +-
21 files changed, 581 insertions(+), 264 deletions(-)
diff --git a/.github/workflows/go-ci.yml b/.github/workflows/go-ci.yml
index b17433228..b026b2614 100644
--- a/.github/workflows/go-ci.yml
+++ b/.github/workflows/go-ci.yml
@@ -93,6 +93,9 @@ jobs:
- name: Run Arrow ownership assertions
if: matrix.os == 'ubuntu-latest' && matrix.go == '1.27.1'
run: make test-assert
+ - name: Vet (including integration-tagged files)
+ if: matrix.os == 'ubuntu-latest' && matrix.go == '1.27.1'
+ run: make vet
- name: Save Go module and build cache
if: github.ref == 'refs/heads/main'
uses: actions/cache/save@55cc8345863c7cc4c66a329aec7e433d2d1c52a9 #
v6.1.0
diff --git a/Makefile b/Makefile
index 2e49d69be..162278efb 100644
--- a/Makefile
+++ b/Makefile
@@ -17,7 +17,7 @@
# golangci-lint version (keep in sync with CI and README)
GOLANGCI_LINT_VERSION := v2.14.0
-.PHONY: test test-assert test-race lint lint-install integration-setup
integration-setup-spark4 integration-test integration-scanner integration-io
integration-rest integration-rest-scan-planning integration-spark
integration-hadoop integration-down integration-logs docs-gen
+.PHONY: test test-assert test-race vet lint lint-install integration-setup
integration-setup-spark4 integration-test integration-scanner integration-io
integration-rest integration-rest-scan-planning integration-spark
integration-hadoop integration-down integration-logs docs-gen
test:
go test -v ./...
@@ -28,6 +28,13 @@ test-assert:
test-race:
go test -race -v ./...
+# vet also type-checks the //go:build integration files, which the default
+# toolchain skips. Without the tagged pass, code behind that tag (for example
+# table/maintenance's integration test) can break without any target noticing.
+vet:
+ go vet ./...
+ go vet -tags=integration ./...
+
docs-gen:
go run ./website/gen
diff --git a/catalog/glue/glue.go b/catalog/glue/glue.go
index 4f085b34d..7824a601f 100644
--- a/catalog/glue/glue.go
+++ b/catalog/glue/glue.go
@@ -35,6 +35,7 @@ import (
"github.com/apache/iceberg-go/io"
"github.com/apache/iceberg-go/metrics"
"github.com/apache/iceberg-go/table"
+ "github.com/apache/iceberg-go/table/maintenance"
"github.com/apache/iceberg-go/utils"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
@@ -660,7 +661,7 @@ func (c *Catalog) PurgeTable(ctx context.Context,
identifier table.Identifier) e
}
// Physically delete all table files on storage best-effort
- if purgeErr := tbl.PurgeFiles(ctx); purgeErr != nil {
+ if purgeErr := maintenance.PurgeFiles(ctx, tbl); purgeErr != nil {
log.Printf("WARNING: dropped table %s but failed to purge
files: %v", identifier, purgeErr)
}
diff --git a/catalog/hadoop/hadoop.go b/catalog/hadoop/hadoop.go
index 76b93a55a..ec25158f8 100644
--- a/catalog/hadoop/hadoop.go
+++ b/catalog/hadoop/hadoop.go
@@ -45,6 +45,7 @@ import (
"github.com/apache/iceberg-go/io"
"github.com/apache/iceberg-go/metrics"
"github.com/apache/iceberg-go/table"
+ "github.com/apache/iceberg-go/table/maintenance"
"github.com/google/uuid"
)
@@ -915,7 +916,7 @@ func (c *Catalog) PurgeTable(ctx context.Context,
identifier table.Identifier) e
}
// For Hadoop catalog, physical files walk must run BEFORE deleting the
table directory root
- if purgeErr := tbl.PurgeFiles(ctx); purgeErr != nil {
+ if purgeErr := maintenance.PurgeFiles(ctx, tbl); purgeErr != nil {
log.Printf("WARNING: failing to purge some files in Hadoop
table %s: %v", identifier, purgeErr)
}
diff --git a/catalog/hive/hive.go b/catalog/hive/hive.go
index 8bcf9e77f..30670bd96 100644
--- a/catalog/hive/hive.go
+++ b/catalog/hive/hive.go
@@ -32,6 +32,7 @@ import (
"github.com/apache/iceberg-go/io"
"github.com/apache/iceberg-go/metrics"
"github.com/apache/iceberg-go/table"
+ "github.com/apache/iceberg-go/table/maintenance"
"github.com/apache/iceberg-go/view"
"github.com/beltran/gohive/hive_metastore"
)
@@ -466,7 +467,7 @@ func (c *Catalog) PurgeTable(ctx context.Context,
identifier table.Identifier) e
}
// Physically delete all table files on storage best-effort
- if purgeErr := tbl.PurgeFiles(ctx); purgeErr != nil {
+ if purgeErr := maintenance.PurgeFiles(ctx, tbl); purgeErr != nil {
log.Printf("WARNING: dropped table %s but failed to purge
files: %v", identifier, purgeErr)
}
diff --git a/catalog/sql/sql.go b/catalog/sql/sql.go
index 6ca265b1c..0d5f25728 100644
--- a/catalog/sql/sql.go
+++ b/catalog/sql/sql.go
@@ -37,6 +37,7 @@ import (
"github.com/apache/iceberg-go/io"
"github.com/apache/iceberg-go/metrics"
"github.com/apache/iceberg-go/table"
+ "github.com/apache/iceberg-go/table/maintenance"
"github.com/apache/iceberg-go/view"
mysql "github.com/go-sql-driver/mysql"
"github.com/uptrace/bun"
@@ -1101,7 +1102,7 @@ func (c *Catalog) PurgeTable(ctx context.Context,
identifier table.Identifier) e
}
// Physically delete all table files on storage best-effort
- if purgeErr := tbl.PurgeFiles(ctx); purgeErr != nil {
+ if purgeErr := maintenance.PurgeFiles(ctx, tbl); purgeErr != nil {
log.Printf("WARNING: dropped table %s but failed to purge
files: %v", identifier, purgeErr)
}
diff --git a/cmd/iceberg/clean_orphan_files.go
b/cmd/iceberg/clean_orphan_files.go
index 48c46533e..a7fd6e483 100644
--- a/cmd/iceberg/clean_orphan_files.go
+++ b/cmd/iceberg/clean_orphan_files.go
@@ -26,6 +26,7 @@ import (
"github.com/apache/iceberg-go/catalog"
"github.com/apache/iceberg-go/table"
+ "github.com/apache/iceberg-go/table/maintenance"
"github.com/pterm/pterm"
)
@@ -38,22 +39,22 @@ func runCleanOrphanFiles(ctx context.Context, output
Output, cat catalog.Catalog
tbl := loadTable(ctx, output, cat, cmd.TableID)
- opts := []table.OrphanCleanupOption{
- table.WithFilesOlderThan(olderThan),
+ opts := []maintenance.OrphanCleanupOption{
+ maintenance.WithCleanupFilesOlderThan(olderThan),
}
if cmd.Location != "" {
- opts = append(opts, table.WithLocation(cmd.Location))
+ opts = append(opts,
maintenance.WithCleanupLocation(cmd.Location))
}
- plan, err := tbl.PlanOrphanFiles(ctx, opts...)
+ plan, err := maintenance.PlanOrphanFiles(ctx, tbl, opts...)
if err != nil {
output.Error(fmt.Errorf("orphan file scan failed: %w", err))
os.Exit(1)
}
planFiles := plan.Files()
- result := table.OrphanCleanupResult{
+ result := maintenance.OrphanCleanupResult{
OrphanFileLocations: planFiles,
OrphanFiles: plan.OrphanFiles(),
TotalSizeBytes: plan.TotalSizeBytes(),
@@ -86,7 +87,7 @@ func runCleanOrphanFiles(ctx context.Context, output Output,
cat catalog.Catalog
os.Exit(1)
}
- deleteResult, err := tbl.ExecuteOrphanCleanup(ctx, plan)
+ deleteResult, err := maintenance.ExecuteOrphanCleanup(ctx, tbl, plan)
if err != nil {
output.Error(fmt.Errorf("orphan file deletion failed: %w", err))
os.Exit(1)
@@ -105,7 +106,7 @@ func shouldPrintCleanOrphanPreview(output Output) bool {
}
}
-func buildCleanOrphanFilesResult(tbl *table.Table, result
table.OrphanCleanupResult, dryRun bool) CleanOrphanFilesResult {
+func buildCleanOrphanFilesResult(tbl *table.Table, result
maintenance.OrphanCleanupResult, dryRun bool) CleanOrphanFilesResult {
var entries []OrphanFileEntry
if dryRun {
entries = make([]OrphanFileEntry, 0, len(result.OrphanFiles))
diff --git a/cmd/iceberg/clean_orphan_files_test.go
b/cmd/iceberg/clean_orphan_files_test.go
index b30cf6b5f..f7ad8b5fb 100644
--- a/cmd/iceberg/clean_orphan_files_test.go
+++ b/cmd/iceberg/clean_orphan_files_test.go
@@ -28,6 +28,7 @@ import (
"github.com/apache/iceberg-go/catalog"
icebergio "github.com/apache/iceberg-go/io"
"github.com/apache/iceberg-go/table"
+ "github.com/apache/iceberg-go/table/maintenance"
"github.com/pterm/pterm"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -127,9 +128,9 @@ func TestBuildCleanOrphanFilesResultDryRun(t *testing.T) {
tbl := table.New([]string{"db", "orphans"}, meta, "", nil, nil)
- orphanResult := table.OrphanCleanupResult{
+ orphanResult := maintenance.OrphanCleanupResult{
OrphanFileLocations: []string{"s3://bucket/data/file1.parquet",
"s3://bucket/data/file2.parquet"},
- OrphanFiles: []table.OrphanFile{
+ OrphanFiles: []maintenance.OrphanFile{
{Path: "s3://bucket/data/file1.parquet", SizeBytes:
1024},
{Path: "s3://bucket/data/file2.parquet", SizeBytes:
2048},
},
@@ -175,9 +176,9 @@ func TestBuildCleanOrphanFilesResultDeleted(t *testing.T) {
tbl := table.New([]string{"db", "orphans"}, meta, "", nil, nil)
- orphanResult := table.OrphanCleanupResult{
+ orphanResult := maintenance.OrphanCleanupResult{
OrphanFileLocations: []string{"s3://bucket/data/file1.parquet"},
- OrphanFiles: []table.OrphanFile{
+ OrphanFiles: []maintenance.OrphanFile{
{Path: "s3://bucket/data/file1.parquet", SizeBytes:
2048},
},
DeletedFiles: []string{"s3://bucket/data/file1.parquet"},
@@ -220,12 +221,12 @@ func
TestBuildCleanOrphanFilesResultDeletedFiltersToDeletedFiles(t *testing.T) {
tbl := table.New([]string{"db", "orphans"}, meta, "", nil, nil)
- orphanResult := table.OrphanCleanupResult{
+ orphanResult := maintenance.OrphanCleanupResult{
OrphanFileLocations: []string{
"s3://bucket/data/file1.parquet",
"s3://bucket/data/file2.parquet",
},
- OrphanFiles: []table.OrphanFile{
+ OrphanFiles: []maintenance.OrphanFile{
{Path: "s3://bucket/data/file1.parquet", SizeBytes:
1024},
{Path: "s3://bucket/data/file2.parquet", SizeBytes:
2048},
},
diff --git a/table/maintenance/manifest_test_helpers_test.go
b/table/maintenance/manifest_test_helpers_test.go
new file mode 100644
index 000000000..86ca1fc1a
--- /dev/null
+++ b/table/maintenance/manifest_test_helpers_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 maintenance
+
+import (
+ "bytes"
+ "fmt"
+ stdfs "io/fs"
+ "sync"
+ "testing"
+ "time"
+
+ "github.com/apache/iceberg-go"
+ iceio "github.com/apache/iceberg-go/io"
+ "github.com/stretchr/testify/require"
+)
+
+// memFile serves manifest and manifest-list bytes staged in trackingIO.
+// A *bytes.Reader already satisfies Read/Seek/ReadAt; memFile adds the
+// Stat/Close that iceio.File additionally requires. Stat returns a nil
+// FileInfo because the manifest readers only ever read the stream.
+type memFile struct {
+ *bytes.Reader
+}
+
+func (memFile) Stat() (stdfs.FileInfo, error) { return nil, nil }
+func (memFile) Close() error { return nil }
+
+// trackingIO is a minimal in-memory IO for the orphan-cleanup reference-set
+// tests. WriteFile stages file bytes in a map and Open serves them back, so a
+// manifest or manifest list can be written once and read back by
+// getReferencedFiles. It implements only the Open/Remove surface iceio.IO
+// requires, plus WriteFile for the manifest helpers.
+type trackingIO struct {
+ files map[string][]byte
+}
+
+func newTrackingIO() *trackingIO {
+ return &trackingIO{files: make(map[string][]byte)}
+}
+
+func (t *trackingIO) Open(name string) (iceio.File, error) {
+ data, ok := t.files[name]
+ if !ok {
+ return nil, stdfs.ErrNotExist
+ }
+
+ return memFile{bytes.NewReader(data)}, nil
+}
+
+func (t *trackingIO) WriteFile(name string, content []byte) error {
+ t.files[name] = append([]byte(nil), content...)
+
+ return nil
+}
+
+func (t *trackingIO) Remove(name string) error {
+ delete(t.files, name)
+
+ return nil
+}
+
+// trackingCallsIO wraps trackingIO to count Open and Remove calls per path.
+type trackingCallsIO struct {
+ *trackingIO
+ mu sync.Mutex
+ openCount map[string]int
+ removeCount map[string]int
+}
+
+func newTrackingCallsIO() *trackingCallsIO {
+ return &trackingCallsIO{
+ trackingIO: newTrackingIO(),
+ openCount: make(map[string]int),
+ removeCount: make(map[string]int),
+ }
+}
+
+func (c *trackingCallsIO) Open(name string) (iceio.File, error) {
+ c.mu.Lock()
+ c.openCount[name]++
+ c.mu.Unlock()
+
+ return c.trackingIO.Open(name)
+}
+
+func (c *trackingCallsIO) Remove(name string) error {
+ c.mu.Lock()
+ c.removeCount[name]++
+ c.mu.Unlock()
+
+ return c.trackingIO.Remove(name)
+}
+
+// writeManifest writes a v2 data manifest with a single ADDED
+// entry pointing at dataPath into tio.files at manifestPath, and returns a
+// ManifestFile descriptor with seqNum pre-assigned so the same descriptor
+// can be referenced from multiple manifest lists.
+func writeManifest(t testing.TB, tio *trackingIO, snapshotID, seqNum int64,
manifestPath, dataPath string) iceberg.ManifestFile {
+ t.Helper()
+
+ dataSchema := iceberg.NewSchema(0,
+ iceberg.NestedField{ID: 1, Name: "x", Type:
iceberg.PrimitiveTypes.Int64, Required: true},
+ )
+ spec := iceberg.NewPartitionSpec()
+
+ df, err := iceberg.NewDataFileBuilder(
+ spec, iceberg.EntryContentData, dataPath, iceberg.ParquetFile,
+ nil, nil, nil, 1, 1024,
+ )
+ require.NoError(t, err)
+
+ entry := iceberg.NewManifestEntryBuilder(iceberg.EntryStatusADDED,
&snapshotID, df.Build()).
+ SequenceNum(seqNum).
+ Build()
+
+ var buf bytes.Buffer
+ _, err = iceberg.WriteManifest(manifestPath, &buf, 2, spec, dataSchema,
snapshotID, []iceberg.ManifestEntry{entry})
+ require.NoError(t, err)
+ require.NoError(t, tio.WriteFile(manifestPath, buf.Bytes()))
+
+ return iceberg.NewManifestFile(2, manifestPath, int64(buf.Len()), 0,
snapshotID).
+ SequenceNum(seqNum, seqNum).
+ AddedFiles(1).
+ AddedRows(1).
+ Build()
+}
+
+// writeManifestList writes a v2 manifest list referencing the
+// given manifests into tio.files at listPath.
+func writeManifestList(t testing.TB, tio *trackingIO, snapshotID int64,
listPath string, manifests []iceberg.ManifestFile) {
+ t.Helper()
+
+ var buf bytes.Buffer
+ seqNum := int64(1)
+ require.NoError(t, iceberg.WriteManifestList(2, &buf, snapshotID, nil,
&seqNum, 0, manifests))
+ require.NoError(t, tio.WriteFile(listPath, buf.Bytes()))
+}
+
+// metaJSONOpts configures the metadata document built by buildMetaJSON.
+type metaJSONOpts struct {
+ snapshots string
+ statistics string
+ partitionStatistics string
+}
+
+// buildMetaJSON returns the minimal v2 metadata document that
ParseMetadataString will accept
+func buildMetaJSON(o metaJSONOpts) string {
+ return fmt.Sprintf(`{
+ "format-version": 2,
+ "table-uuid": "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa",
+ "location": "s3://bucket/table",
+ "last-sequence-number": 0,
+ "last-updated-ms": 1000,
+ "last-column-id": 1,
+ "current-schema-id": 0,
+ "schemas":
[{"type":"struct","schema-id":0,"fields":[{"id":1,"name":"x","required":true,"type":"long"}]}],
+ "default-spec-id": 0,
+ "partition-specs": [{"spec-id":0,"fields":[]}],
+ "last-partition-id": 0,
+ "default-sort-order-id": 0,
+ "sort-orders": [{"order-id":0,"fields":[]}],
+ "snapshots": [%s],
+ "statistics": [%s],
+ "partition-statistics": [%s]
+ }`, o.snapshots, o.statistics, o.partitionStatistics)
+}
+
+// manifestTrackingIO wraps an IO to record the peak number of concurrently
+// open files, so tests can assert that manifest-list reads run concurrently
+// yet stay within the configured worker limit.
+type manifestTrackingIO struct {
+ iceio.IO
+ mu sync.Mutex
+ open int
+ maxOpen int
+ delay time.Duration
+}
+
+func (fs *manifestTrackingIO) Open(name string) (iceio.File, error) {
+ f, err := fs.IO.Open(name)
+ if err != nil {
+ return nil, err
+ }
+
+ fs.mu.Lock()
+ fs.open++
+ if fs.open > fs.maxOpen {
+ fs.maxOpen = fs.open
+ }
+ fs.mu.Unlock()
+ time.Sleep(fs.delay)
+
+ return &manifestTrackingFile{File: f, onClose: func() {
+ fs.mu.Lock()
+ fs.open--
+ fs.mu.Unlock()
+ }}, nil
+}
+
+type manifestTrackingFile struct {
+ iceio.File
+ onClose func()
+ once sync.Once
+}
+
+func (f *manifestTrackingFile) Close() error {
+ f.once.Do(f.onClose)
+
+ return f.File.Close()
+}
+
+// benchmarkDelayIO wraps an IO and adds a fixed per-Open delay, modelling the
+// round trip to a remote object store in the manifest-list benchmark.
+type benchmarkDelayIO struct {
+ iceio.IO
+ delay time.Duration
+}
+
+func (fs *benchmarkDelayIO) Open(name string) (iceio.File, error) {
+ f, err := fs.IO.Open(name)
+ if err != nil {
+ return nil, err
+ }
+ time.Sleep(fs.delay)
+
+ return f, nil
+}
diff --git a/table/orphan_cleanup.go b/table/maintenance/orphan_cleanup.go
similarity index 92%
rename from table/orphan_cleanup.go
rename to table/maintenance/orphan_cleanup.go
index 22d6a6019..bfae0048c 100644
--- a/table/orphan_cleanup.go
+++ b/table/maintenance/orphan_cleanup.go
@@ -15,7 +15,12 @@
// specific language governing permissions and limitations
// under the License.
-package table
+// Package maintenance provides table maintenance operations, such as orphan
+// file cleanup, that operate on a loaded [table.Table]. The operations are
+// free functions that take the table as their first argument, for example:
+//
+// result, err := maintenance.DeleteOrphanFiles(ctx, tbl,
maintenance.WithCleanupFilesOlderThan(72*time.Hour))
+package maintenance
import (
"context"
@@ -37,7 +42,8 @@ import (
"github.com/apache/iceberg-go"
iceberginternal "github.com/apache/iceberg-go/internal"
"github.com/apache/iceberg-go/internal/fileuri"
- iceio "github.com/apache/iceberg-go/io"
+ "github.com/apache/iceberg-go/io"
+ "github.com/apache/iceberg-go/table"
"golang.org/x/sync/errgroup"
)
@@ -93,18 +99,25 @@ type orphanCleanupConfig struct {
const defaultPurgeMaxConcurrency = 32
+// OrphanCleanupOption configures an orphan-cleanup run. Options are passed to
+// DeleteOrphanFiles, PlanOrphanFiles, and ExecuteOrphanCleanup.
type OrphanCleanupOption func(*orphanCleanupConfig)
-func WithLocation(location string) OrphanCleanupOption {
+// WithCleanupLocation sets the location to scan for orphan files. Defaults to
+// the table's data location.
+func WithCleanupLocation(location string) OrphanCleanupOption {
return func(cfg *orphanCleanupConfig) {
- rejectPlanOption(cfg, "WithLocation")
+ rejectPlanOption(cfg, "WithCleanupLocation")
cfg.location = location
}
}
-func WithFilesOlderThan(duration time.Duration) OrphanCleanupOption {
+// WithCleanupFilesOlderThan only considers files last modified before now
minus
+// duration, so recently written files are never treated as orphans. Defaults
to
+// 72h. The duration must be non-negative.
+func WithCleanupFilesOlderThan(duration time.Duration) OrphanCleanupOption {
return func(cfg *orphanCleanupConfig) {
- rejectPlanOption(cfg, "WithFilesOlderThan")
+ rejectPlanOption(cfg, "WithCleanupFilesOlderThan")
cfg.olderThan = duration
if duration < 0 && cfg.validationErr == nil {
cfg.validationErr = errors.New("orphan cleanup age must
be non-negative")
@@ -112,15 +125,17 @@ func WithFilesOlderThan(duration time.Duration)
OrphanCleanupOption {
}
}
-func WithDryRun(enabled bool) OrphanCleanupOption {
+// WithCleanupDryRun reports which files would be deleted without deleting
them.
+// Defaults to false.
+func WithCleanupDryRun(enabled bool) OrphanCleanupOption {
return func(cfg *orphanCleanupConfig) {
cfg.dryRun = enabled
}
}
-// WithDeleteFunc sets a custom delete function. If not provided, the table's
FileIO
+// WithCleanupDeleteFunc sets a custom delete function. If not provided, the
table's FileIO
// delete method will be used.
-func WithDeleteFunc(deleteFunc func(string) error) OrphanCleanupOption {
+func WithCleanupDeleteFunc(deleteFunc func(string) error) OrphanCleanupOption {
return func(cfg *orphanCleanupConfig) {
cfg.deleteFunc = deleteFunc
}
@@ -307,14 +322,14 @@ func rejectPlanOption(cfg *orphanCleanupConfig, name
string) {
// DeleteOrphanFiles identifies files under a table location that are no longer
// referenced by table metadata and deletes them unless dry-run is enabled.
//
-// The table filesystem must implement iceio.ListableIO so orphan cleanup can
+// The table filesystem must implement io.ListableIO so orphan cleanup can
// fully enumerate candidate files before deciding what is safe to delete.
-func (t Table) DeleteOrphanFiles(ctx context.Context, opts
...OrphanCleanupOption) (OrphanCleanupResult, error) {
+func DeleteOrphanFiles(ctx context.Context, tbl *table.Table, opts
...OrphanCleanupOption) (OrphanCleanupResult, error) {
cfg := newOrphanCleanupConfig(opts...)
if cfg.validationErr != nil {
return OrphanCleanupResult{}, cfg.validationErr
}
- plan, err := t.planOrphanFiles(ctx, cfg)
+ plan, err := planOrphanFiles(ctx, tbl, cfg)
if err != nil {
return OrphanCleanupResult{}, err
}
@@ -322,24 +337,24 @@ func (t Table) DeleteOrphanFiles(ctx context.Context,
opts ...OrphanCleanupOptio
return plan.result(), nil
}
- return t.executeOrphanCleanup(ctx, plan, cfg)
+ return executeOrphanCleanup(ctx, tbl, plan, cfg)
}
// PlanOrphanFiles identifies orphan files without deleting anything. The
// returned plan can be shown to a user and passed to ExecuteOrphanCleanup to
// delete exactly that set.
-func (t Table) PlanOrphanFiles(ctx context.Context, opts
...OrphanCleanupOption) (OrphanCleanupPlan, error) {
+func PlanOrphanFiles(ctx context.Context, tbl *table.Table, opts
...OrphanCleanupOption) (OrphanCleanupPlan, error) {
cfg := newOrphanCleanupConfig(opts...)
if cfg.validationErr != nil {
return OrphanCleanupPlan{}, cfg.validationErr
}
- return t.planOrphanFiles(ctx, cfg)
+ return planOrphanFiles(ctx, tbl, cfg)
}
// ExecuteOrphanCleanup deletes exactly the files in plan. It does not perform
// another orphan scan, so files appearing after planning are not included.
-func (t Table) ExecuteOrphanCleanup(ctx context.Context, plan
OrphanCleanupPlan, opts ...OrphanCleanupOption) (OrphanCleanupResult, error) {
+func ExecuteOrphanCleanup(ctx context.Context, tbl *table.Table, plan
OrphanCleanupPlan, opts ...OrphanCleanupOption) (OrphanCleanupResult, error) {
cfg := newExecutionOrphanCleanupConfig(opts...)
if cfg.validationErr != nil {
return OrphanCleanupResult{}, cfg.validationErr
@@ -351,7 +366,7 @@ func (t Table) ExecuteOrphanCleanup(ctx context.Context,
plan OrphanCleanupPlan,
return plan.result(), nil
}
- return t.executeOrphanCleanup(ctx, plan, cfg)
+ return executeOrphanCleanup(ctx, tbl, plan, cfg)
}
type scannedFile struct {
@@ -364,15 +379,15 @@ type referencedFileIndex struct {
byPath map[string][]string
}
-func (t Table) planOrphanFiles(ctx context.Context, cfg *orphanCleanupConfig)
(OrphanCleanupPlan, error) {
- fs, err := t.fsF(ctx)
+func planOrphanFiles(ctx context.Context, tbl *table.Table, cfg
*orphanCleanupConfig) (OrphanCleanupPlan, error) {
+ fs, err := tbl.FS(ctx)
if err != nil {
return OrphanCleanupPlan{}, fmt.Errorf("failed to get
filesystem: %w", err)
}
scanLocation := cfg.location
if scanLocation == "" {
- scanLocation = t.metadata.Location()
+ scanLocation = tbl.Metadata().Location()
}
// Run the S3 walk and referenced-file collection concurrently.
@@ -384,7 +399,7 @@ func (t Table) planOrphanFiles(ctx context.Context, cfg
*orphanCleanupConfig) (O
g, gctx := errgroup.WithContext(ctx)
g.Go(func() error {
var err error
- referencedFiles, err = t.getReferencedFiles(gctx, fs,
cfg.maxConcurrency, true)
+ referencedFiles, err = getReferencedFiles(gctx, tbl, fs,
cfg.maxConcurrency, true)
return err
})
@@ -436,8 +451,8 @@ func (t Table) planOrphanFiles(ctx context.Context, cfg
*orphanCleanupConfig) (O
}, nil
}
-func (t Table) executeOrphanCleanup(ctx context.Context, plan
OrphanCleanupPlan, cfg *orphanCleanupConfig) (OrphanCleanupResult, error) {
- fs, err := t.fsF(ctx)
+func executeOrphanCleanup(ctx context.Context, tbl *table.Table, plan
OrphanCleanupPlan, cfg *orphanCleanupConfig) (OrphanCleanupResult, error) {
+ fs, err := tbl.FS(ctx)
if err != nil {
return OrphanCleanupResult{}, fmt.Errorf("failed to get
filesystem: %w", err)
}
@@ -470,14 +485,14 @@ func (t Table) executeOrphanCleanup(ctx context.Context,
plan OrphanCleanupPlan,
// comparison identity before normalization discards information. The bool
value
// distinguishes data files (true) from metadata files (false), which is used
by
// PurgeFiles to respect gc.enabled.
-func (t Table) getReferencedFiles(ctx context.Context, fs iceio.IO,
maxConcurrency int, discardDeleted bool) (map[string]bool, error) {
+func getReferencedFiles(ctx context.Context, tbl *table.Table, fs io.IO,
maxConcurrency int, discardDeleted bool) (map[string]bool, error) {
referenced := make(map[string]bool)
- metadata := t.metadata
+ metadata := tbl.Metadata()
for entry := range metadata.PreviousFiles() {
referenced[entry.MetadataFile] = false
}
- referenced[t.metadataLocation] = false
+ referenced[tbl.MetadataLocation()] = false
// Add version hint file (for Hadoop-style tables)
// Following Java's ReachableFileUtil.versionHintLocation() logic:
@@ -611,8 +626,8 @@ func (t Table) getReferencedFiles(ctx context.Context, fs
iceio.IO, maxConcurren
return referenced, nil
}
-func walkDirectory(fsys iceio.IO, root string, fn func(path string, info
stdfs.FileInfo) error) error {
- if listable, ok := fsys.(iceio.ListableIO); ok {
+func walkDirectory(fsys io.IO, root string, fn func(path string, info
stdfs.FileInfo) error) error {
+ if listable, ok := fsys.(io.ListableIO); ok {
return listable.WalkDir(root, func(path string, d
stdfs.DirEntry, err error) error {
if err != nil {
return err
@@ -631,7 +646,7 @@ func walkDirectory(fsys iceio.IO, root string, fn func(path
string, info stdfs.F
})
}
- return fmt.Errorf("filesystem %T does not implement iceio.ListableIO",
fsys)
+ return fmt.Errorf("filesystem %T does not implement io.ListableIO",
fsys)
}
func isFileOrphan(
@@ -702,14 +717,14 @@ func newReferencedFileIndex(referencedFiles
map[string]bool, cfg *orphanCleanupC
return index
}
-func deleteFiles(ctx context.Context, fs iceio.IO, orphanFiles []string, cfg
*orphanCleanupConfig) ([]string, error) {
+func deleteFiles(ctx context.Context, fs io.IO, orphanFiles []string, cfg
*orphanCleanupConfig) ([]string, error) {
if len(orphanFiles) == 0 {
return nil, nil
}
// Use bulk delete when available and no custom deleteFunc is set.
if cfg.deleteFunc == nil {
- if bulk, ok := fs.(iceio.BulkRemovableIO); ok {
+ if bulk, ok := fs.(io.BulkRemovableIO); ok {
return bulk.DeleteFiles(ctx, orphanFiles)
}
}
@@ -734,7 +749,7 @@ func deleteFiles(ctx context.Context, fs iceio.IO,
orphanFiles []string, cfg *or
)
}
-func deleteFilesSequential(ctx context.Context, fs iceio.IO, orphanFiles
[]string, cfg *orphanCleanupConfig) ([]string, error) {
+func deleteFilesSequential(ctx context.Context, fs io.IO, orphanFiles
[]string, cfg *orphanCleanupConfig) ([]string, error) {
var deletedFiles []string
deleteFunc := fs.Remove
@@ -1324,23 +1339,23 @@ func pathPrefix(path string) (scheme, authority string,
ok bool) {
// does not get out of sync with storage. Non-bulk deletion invokes Remove
// concurrently and is bounded to defaultPurgeMaxConcurrency operations because
// it is typically I/O-bound.
-func (t Table) PurgeFiles(ctx context.Context) error {
- gcEnabled := isGCEnabled(t.Metadata().Properties())
+func PurgeFiles(ctx context.Context, tbl *table.Table) error {
+ gcEnabled := table.IsGCEnabled(tbl.Properties())
- fs, err := t.FS(ctx)
+ fs, err := tbl.FS(ctx)
if err != nil {
return fmt.Errorf("failed to load filesystem for table purge:
%w", err)
}
var errs []error
fileSet := make(map[string]string)
- location := t.metadata.Location()
+ location := tbl.Metadata().Location()
// 1. Walk the table location directory tree to capture all local files
// Only walk the directory if gc.enabled=true to prevent accidental
deletion
// of unreferenced branched data files.
if gcEnabled {
- if listable, ok := fs.(iceio.ListableIO); ok {
+ if listable, ok := fs.(io.ListableIO); ok {
walkErr := listable.WalkDir(location, func(path string,
d stdfs.DirEntry, err error) error {
if err := ctx.Err(); err != nil {
return err
@@ -1365,7 +1380,7 @@ func (t Table) PurgeFiles(ctx context.Context) error {
}
// 2. Union in manifest-referenced and metadata files (which might be
outside the table location)
- referencedFiles, refErr := t.getReferencedFiles(ctx, fs,
runtime.GOMAXPROCS(0), false)
+ referencedFiles, refErr := getReferencedFiles(ctx, tbl, fs,
runtime.GOMAXPROCS(0), false)
if refErr != nil {
return fmt.Errorf("failed to get referenced files: %w", refErr)
}
@@ -1391,7 +1406,7 @@ func (t Table) PurgeFiles(ctx context.Context) error {
slices.Sort(files)
if len(files) > 0 {
- if bulk, ok := fs.(iceio.BulkRemovableIO); ok {
+ if bulk, ok := fs.(io.BulkRemovableIO); ok {
_, bulkErr := bulk.DeleteFiles(ctx, files)
if bulkErr != nil {
errs = append(errs, fmt.Errorf("bulk deletion
failed: %w", bulkErr))
diff --git a/table/orphan_cleanup_bench_test.go
b/table/maintenance/orphan_cleanup_bench_test.go
similarity index 89%
rename from table/orphan_cleanup_bench_test.go
rename to table/maintenance/orphan_cleanup_bench_test.go
index 8e7dd5d4c..2f8aee49f 100644
--- a/table/orphan_cleanup_bench_test.go
+++ b/table/maintenance/orphan_cleanup_bench_test.go
@@ -15,7 +15,7 @@
// specific language governing permissions and limitations
// under the License.
-package table
+package maintenance
import (
"context"
@@ -25,7 +25,7 @@ import (
"time"
"github.com/apache/iceberg-go"
- iceio "github.com/apache/iceberg-go/io"
+ "github.com/apache/iceberg-go/table"
)
var orphanCleanupBenchmarkSink string
@@ -86,7 +86,7 @@ func BenchmarkGetReferencedFilesManifestLists(b *testing.B) {
i, i*1000, listPath))
}
- meta, err :=
ParseMetadataString(buildMetaJSON(metaJSONOpts{
+ meta, err :=
table.ParseMetadataString(buildMetaJSON(metaJSONOpts{
snapshots: strings.Join(snapshotJSON,
","),
}))
if err != nil {
@@ -94,14 +94,14 @@ func BenchmarkGetReferencedFilesManifestLists(b *testing.B)
{
}
fs := &benchmarkDelayIO{IO: baseIO, delay:
time.Millisecond}
- tbl := New(Identifier{"ns",
"orphan-cleanup-benchmark"}, meta,
+ tbl := table.New(table.Identifier{"ns",
"orphan-cleanup-benchmark"}, meta,
"metadata.json", testFSF(fs), nil)
b.ReportAllocs()
b.ReportMetric(float64(snapshotCount),
"manifest_lists/op")
b.ResetTimer()
for b.Loop() {
- if _, err :=
tbl.getReferencedFiles(context.Background(), fs, maxWorkers, true); err != nil {
+ if _, err :=
getReferencedFiles(context.Background(), tbl, fs, maxWorkers, true); err != nil
{
b.Fatal(err)
}
}
@@ -110,21 +110,6 @@ func BenchmarkGetReferencedFilesManifestLists(b
*testing.B) {
}
}
-type benchmarkDelayIO struct {
- iceio.IO
- delay time.Duration
-}
-
-func (fs *benchmarkDelayIO) Open(name string) (iceio.File, error) {
- f, err := fs.IO.Open(name)
- if err != nil {
- return nil, err
- }
- time.Sleep(fs.delay)
-
- return f, nil
-}
-
// 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
diff --git a/table/orphan_cleanup_integration_test.go
b/table/maintenance/orphan_cleanup_integration_test.go
similarity index 91%
rename from table/orphan_cleanup_integration_test.go
rename to table/maintenance/orphan_cleanup_integration_test.go
index 0b8e0dc38..b72358461 100644
--- a/table/orphan_cleanup_integration_test.go
+++ b/table/maintenance/orphan_cleanup_integration_test.go
@@ -17,7 +17,7 @@
//go:build integration
-package table_test
+package maintenance_test
import (
"context"
@@ -36,6 +36,7 @@ import (
"github.com/apache/iceberg-go/io"
_ "github.com/apache/iceberg-go/io/gocloud"
"github.com/apache/iceberg-go/table"
+ "github.com/apache/iceberg-go/table/maintenance"
"github.com/stretchr/testify/suite"
)
@@ -220,10 +221,10 @@ func (s *OrphanCleanupIntegrationSuite)
TestOrphanCleanupDryRun() {
// Run dry-run orphan cleanup (scan table root where both data and
metadata exist)
s.T().Logf("Scanning location: %s", tbl.Location())
- result, err := tbl.DeleteOrphanFiles(s.ctx,
- table.WithDryRun(true),
- table.WithLocation(tbl.Location()), // Scan the table root
location
- table.WithFilesOlderThan(0), // Consider files created
before the scan
+ result, err := maintenance.DeleteOrphanFiles(s.ctx, tbl,
+ maintenance.WithCleanupDryRun(true),
+ maintenance.WithCleanupLocation(tbl.Location()), // Scan the
table root location
+ maintenance.WithCleanupFilesOlderThan(0), // Consider
files created before the scan
)
s.Require().NoError(err)
@@ -268,9 +269,9 @@ func (s *OrphanCleanupIntegrationSuite)
TestOrphanCleanupActualDeletion() {
}
// Run actual orphan cleanup
- result, err := tbl.DeleteOrphanFiles(s.ctx,
- table.WithDryRun(false),
- table.WithFilesOlderThan(0), // Consider files created before
the scan
+ result, err := maintenance.DeleteOrphanFiles(s.ctx, tbl,
+ maintenance.WithCleanupDryRun(false),
+ maintenance.WithCleanupFilesOlderThan(0), // Consider files
created before the scan
)
s.Require().NoError(err)
@@ -309,10 +310,10 @@ func (s *OrphanCleanupIntegrationSuite)
TestOrphanCleanupCustomLocation() {
time.Sleep(100 * time.Millisecond)
// Run orphan cleanup on table root (which includes the custom
subdirectory)
- result, err := tbl.DeleteOrphanFiles(s.ctx,
- table.WithLocation(tbl.Location()), // Scan table root to find
all subdirectories
- table.WithFilesOlderThan(0),
- table.WithDryRun(false),
+ result, err := maintenance.DeleteOrphanFiles(s.ctx, tbl,
+ maintenance.WithCleanupLocation(tbl.Location()), // Scan table
root to find all subdirectories
+ maintenance.WithCleanupFilesOlderThan(0),
+ maintenance.WithCleanupDryRun(false),
)
s.Require().NoError(err)
@@ -358,10 +359,10 @@ func (s *OrphanCleanupIntegrationSuite)
TestOrphanCleanupWithConcurrency() {
time.Sleep(100 * time.Millisecond)
// Run orphan cleanup with specific concurrency
- result, err := tbl.DeleteOrphanFiles(s.ctx,
- table.WithCleanupMaxConcurrency(tc.concurrency),
- table.WithFilesOlderThan(0),
- table.WithDryRun(false),
+ result, err := maintenance.DeleteOrphanFiles(s.ctx, tbl,
+
maintenance.WithCleanupMaxConcurrency(tc.concurrency),
+ maintenance.WithCleanupFilesOlderThan(0),
+ maintenance.WithCleanupDryRun(false),
)
s.Require().NoError(err)
@@ -398,11 +399,11 @@ func (s *OrphanCleanupIntegrationSuite)
TestOrphanCleanupCustomDeleteFunction()
return fs.Remove(filePath)
}
- result, err := tbl.DeleteOrphanFiles(s.ctx,
- table.WithDeleteFunc(customDeleteFunc),
- table.WithFilesOlderThan(0),
- table.WithDryRun(false),
- table.WithDryRun(false),
+ result, err := maintenance.DeleteOrphanFiles(s.ctx, tbl,
+ maintenance.WithCleanupDeleteFunc(customDeleteFunc),
+ maintenance.WithCleanupFilesOlderThan(0),
+ maintenance.WithCleanupDryRun(false),
+ maintenance.WithCleanupDryRun(false),
)
s.Require().NoError(err)
diff --git a/table/orphan_cleanup_test.go
b/table/maintenance/orphan_cleanup_test.go
similarity index 90%
rename from table/orphan_cleanup_test.go
rename to table/maintenance/orphan_cleanup_test.go
index 962ff670e..6709ad95d 100644
--- a/table/orphan_cleanup_test.go
+++ b/table/maintenance/orphan_cleanup_test.go
@@ -15,7 +15,7 @@
// specific language governing permissions and limitations
// under the License.
-package table
+package maintenance
import (
"context"
@@ -36,7 +36,8 @@ import (
"github.com/apache/arrow-go/v18/arrow/array"
"github.com/apache/arrow-go/v18/arrow/memory"
"github.com/apache/iceberg-go"
- "github.com/apache/iceberg-go/io"
+ iceio "github.com/apache/iceberg-go/io"
+ "github.com/apache/iceberg-go/table"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
@@ -62,18 +63,18 @@ func TestPrefixMismatchMode_String(t *testing.T) {
func TestOrphanCleanupOptions(t *testing.T) {
cfg := &orphanCleanupConfig{}
- WithLocation("/test/location")(cfg)
+ WithCleanupLocation("/test/location")(cfg)
assert.Equal(t, "/test/location", cfg.location)
testDuration := 24 * time.Hour
- WithFilesOlderThan(testDuration)(cfg)
+ WithCleanupFilesOlderThan(testDuration)(cfg)
assert.Equal(t, testDuration, cfg.olderThan)
- WithDryRun(true)(cfg)
+ WithCleanupDryRun(true)(cfg)
assert.True(t, cfg.dryRun)
deleteFunc := func(string) error { return nil }
- WithDeleteFunc(deleteFunc)(cfg)
+ WithCleanupDeleteFunc(deleteFunc)(cfg)
assert.NotNil(t, cfg.deleteFunc)
WithCleanupMaxConcurrency(8)(cfg)
@@ -165,11 +166,11 @@ func TestNewOrphanCleanupConfigFlattensURIEquivalences(t
*testing.T) {
func TestOrphanCleanupPlanDoesNotExpandAfterPlanning(t *testing.T) {
ctx := context.Background()
- fs := io.NewMemFS()
+ fs := iceio.NewMemFS()
location := "mem://plan-race/table"
metadataLocation := location + "/metadata/v1.metadata.json"
- meta, err := NewMetadata(iceberg.NewSchema(0), nil, UnsortedSortOrder,
location, nil)
+ meta, err := table.NewMetadata(iceberg.NewSchema(0), nil,
table.UnsortedSortOrder, location, nil)
require.NoError(t, err)
require.NoError(t, fs.WriteFile(metadataLocation, nil))
@@ -177,25 +178,25 @@ func TestOrphanCleanupPlanDoesNotExpandAfterPlanning(t
*testing.T) {
newOrphan := location + "/data/appeared-after-confirmation.parquet"
require.NoError(t, fs.WriteFile(plannedOrphan, []byte("planned")))
- tbl := New(
+ tbl := table.New(
[]string{"db", "plan_race"},
meta,
metadataLocation,
- func(context.Context) (io.IO, error) { return fs, nil },
+ func(context.Context) (iceio.IO, error) { return fs, nil },
nil,
)
- plan, err := tbl.PlanOrphanFiles(ctx, WithFilesOlderThan(time.Hour))
+ plan, err := PlanOrphanFiles(ctx, tbl,
WithCleanupFilesOlderThan(time.Hour))
require.NoError(t, err)
assert.Equal(t, []string{plannedOrphan}, plan.Files())
assert.Equal(t, []OrphanFile{{Path: plannedOrphan, SizeBytes:
int64(len("planned"))}}, plan.OrphanFiles())
assert.False(t, plan.Cutoff().IsZero())
require.NoError(t, fs.WriteFile(newOrphan, []byte("new")))
- _, err = tbl.ExecuteOrphanCleanup(ctx, plan,
WithFilesOlderThan(time.Hour))
- require.ErrorContains(t, err, "WithFilesOlderThan")
+ _, err = ExecuteOrphanCleanup(ctx, tbl, plan,
WithCleanupFilesOlderThan(time.Hour))
+ require.ErrorContains(t, err, "WithCleanupFilesOlderThan")
- result, err := tbl.ExecuteOrphanCleanup(ctx, plan,
WithCleanupMaxConcurrency(1))
+ result, err := ExecuteOrphanCleanup(ctx, tbl, plan,
WithCleanupMaxConcurrency(1))
require.NoError(t, err)
assert.Equal(t, []string{plannedOrphan}, result.DeletedFiles)
@@ -211,8 +212,8 @@ func TestExecuteOrphanCleanupRejectsPlanningOptions(t
*testing.T) {
name string
opt OrphanCleanupOption
}{
- {name: "location", opt: WithLocation("mem://other")},
- {name: "age", opt: WithFilesOlderThan(time.Hour)},
+ {name: "location", opt: WithCleanupLocation("mem://other")},
+ {name: "age", opt: WithCleanupFilesOlderThan(time.Hour)},
{name: "prefix mismatch mode", opt:
WithPrefixMismatchMode(PrefixMismatchIgnore)},
{name: "equal schemes", opt:
WithEqualSchemes(map[string]string{"s3a": "s3"})},
{name: "equal authorities", opt:
WithEqualAuthorities(map[string]string{"old": "new"})},
@@ -220,7 +221,7 @@ func TestExecuteOrphanCleanupRejectsPlanningOptions(t
*testing.T) {
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
- _, err :=
(Table{}).ExecuteOrphanCleanup(context.Background(), OrphanCleanupPlan{},
tt.opt)
+ _, err := ExecuteOrphanCleanup(context.Background(),
&table.Table{}, OrphanCleanupPlan{}, tt.opt)
require.Error(t, err)
assert.Contains(t, err.Error(), "only valid while
planning")
})
@@ -235,7 +236,7 @@ func TestPlanOrphanFilesHonorsModificationTimes(t
*testing.T) {
recentPath := location + "/data/recent.parquet"
now := time.Now()
- meta, err := NewMetadata(iceberg.NewSchema(0), nil, UnsortedSortOrder,
location, nil)
+ meta, err := table.NewMetadata(iceberg.NewSchema(0), nil,
table.UnsortedSortOrder, location, nil)
require.NoError(t, err)
mockFS := &mockListableIO{
entries: []mockWalkEntry{
@@ -245,15 +246,15 @@ func TestPlanOrphanFilesHonorsModificationTimes(t
*testing.T) {
{path: recentPath, info: mockFileInfo{name:
"recent.parquet", size: 20, modTime: now.Add(-10 * time.Minute)}},
},
}
- tbl := New(
- Identifier{"db", "mtime"},
+ tbl := table.New(
+ table.Identifier{"db", "mtime"},
meta,
metadataLocation,
- func(context.Context) (io.IO, error) { return mockFS, nil },
+ func(context.Context) (iceio.IO, error) { return mockFS, nil },
nil,
)
- plan, err := tbl.PlanOrphanFiles(ctx, WithFilesOlderThan(time.Hour))
+ plan, err := PlanOrphanFiles(ctx, tbl,
WithCleanupFilesOlderThan(time.Hour))
require.NoError(t, err)
assert.Equal(t, []string{oldPath}, plan.Files())
assert.Equal(t, []OrphanFile{{Path: oldPath, SizeBytes: 10}},
plan.OrphanFiles())
@@ -287,14 +288,14 @@ func
TestPlanOrphanFilesPreservesPOSIXInterpretationOfAmbiguousUNCReference(t *t
},
} {
t.Run(tt.name, func(t *testing.T) {
- meta, err := NewMetadata(iceberg.NewSchema(0), nil,
UnsortedSortOrder, tt.tableLocation, nil)
+ meta, err := table.NewMetadata(iceberg.NewSchema(0),
nil, table.UnsortedSortOrder, tt.tableLocation, nil)
require.NoError(t, err)
- builder, err := MetadataBuilderFromBase(meta,
tt.metadataPath)
+ builder, err := table.MetadataBuilderFromBase(meta,
tt.metadataPath)
require.NoError(t, err)
- require.NoError(t, builder.SetStatistics(StatisticsFile{
+ require.NoError(t,
builder.SetStatistics(table.StatisticsFile{
SnapshotID: 1,
StatisticsPath: tt.referencedPath,
- BlobMetadata: []BlobMetadata{},
+ BlobMetadata: []table.BlobMetadata{},
FileSizeInBytes: 4,
}))
meta, err = builder.Build()
@@ -304,15 +305,15 @@ func
TestPlanOrphanFilesPreservesPOSIXInterpretationOfAmbiguousUNCReference(t *t
path: tt.listedPath,
info: mockFileInfo{name: "file.parquet", size:
4},
}}}
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
tt.metadataPath,
- func(context.Context) (io.IO, error) { return
fsys, nil },
+ func(context.Context) (iceio.IO, error) {
return fsys, nil },
nil,
)
- plan, err := tbl.PlanOrphanFiles(context.Background(),
WithFilesOlderThan(0))
+ plan, err := PlanOrphanFiles(context.Background(), tbl,
WithCleanupFilesOlderThan(0))
require.NoError(t, err)
assert.Empty(t, plan.Files())
})
@@ -335,14 +336,14 @@ func
TestPlanOrphanFilesMatchesNativeHierarchicalFileURIPath(t *testing.T) {
"file://localhost/C:/Warehouse/Table/Data/file.parquet",
} {
t.Run(reference, func(t *testing.T) {
- meta, err := NewMetadata(iceberg.NewSchema(0), nil,
UnsortedSortOrder, tableLocation, nil)
+ meta, err := table.NewMetadata(iceberg.NewSchema(0),
nil, table.UnsortedSortOrder, tableLocation, nil)
require.NoError(t, err)
- builder, err := MetadataBuilderFromBase(meta,
metadataPath)
+ builder, err := table.MetadataBuilderFromBase(meta,
metadataPath)
require.NoError(t, err)
- require.NoError(t, builder.SetStatistics(StatisticsFile{
+ require.NoError(t,
builder.SetStatistics(table.StatisticsFile{
SnapshotID: 1,
StatisticsPath: reference,
- BlobMetadata: []BlobMetadata{},
+ BlobMetadata: []table.BlobMetadata{},
FileSizeInBytes: 4,
}))
meta, err = builder.Build()
@@ -352,15 +353,15 @@ func
TestPlanOrphanFilesMatchesNativeHierarchicalFileURIPath(t *testing.T) {
path: listedPath,
info: mockFileInfo{name: "file.parquet", size:
4},
}}}
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
metadataPath,
- func(context.Context) (io.IO, error) { return
fsys, nil },
+ func(context.Context) (iceio.IO, error) {
return fsys, nil },
nil,
)
- plan, err := tbl.PlanOrphanFiles(context.Background(),
WithFilesOlderThan(0))
+ plan, err := PlanOrphanFiles(context.Background(), tbl,
WithCleanupFilesOlderThan(0))
require.NoError(t, err)
assert.Empty(t, plan.Files())
})
@@ -1488,21 +1489,24 @@ func TestGetReferencedFiles_IncludesStatisticsFiles(t
*testing.T) {
]
}`
- meta, err := ParseMetadataString(metaJSON)
+ meta, err := table.ParseMetadataString(metaJSON)
require.NoError(t, err)
- tbl := Table{
- metadata: meta,
- metadataLocation:
"s3://bucket/test/location/metadata/v1.metadata.json",
- }
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
+ meta,
+ "s3://bucket/test/location/metadata/v1.metadata.json",
+ nil,
+ nil,
+ )
// No snapshots: FileIO is not used; statistics paths must still be
referenced.
- refs, err := tbl.getReferencedFiles(context.Background(), nil, 1, true)
+ refs, err := getReferencedFiles(context.Background(), tbl, nil, 1, true)
require.NoError(t, err)
assert.Contains(t, refs,
normalizeFilePath("s3://bucket/stats/table-stats.puffin"))
assert.Contains(t, refs,
normalizeFilePath("s3://bucket/stats/part-stats.puffin"))
- assert.Contains(t, refs, normalizeFilePath(tbl.metadataLocation))
+ assert.Contains(t, refs, normalizeFilePath(tbl.MetadataLocation()))
assert.Contains(t, refs,
normalizeFilePath("s3://bucket/test/location/metadata/version-hint.text"))
assert.NotContains(t, refs,
normalizeFilePath("s3:/bucket/test/location/metadata/version-hint.text"))
assert.NotContains(t, refs,
normalizeFilePath("s3://bucket/stats/not-referenced.puffin"))
@@ -1515,7 +1519,7 @@ type mockBulkRemovableIO struct {
bulkPaths []string
}
-func (m *mockBulkRemovableIO) Open(string) (io.File, error) {
+func (m *mockBulkRemovableIO) Open(string) (iceio.File, error) {
return nil, errors.New("not implemented")
}
@@ -1575,7 +1579,7 @@ type mockPlainIO struct {
removed []string
}
-func (m *mockPlainIO) Open(string) (io.File, error) {
+func (m *mockPlainIO) Open(string) (iceio.File, error) {
return nil, errors.New("not implemented")
}
@@ -1874,16 +1878,16 @@ func TestPurgeFilesDeletesNonBulkFilesConcurrently(t
*testing.T) {
release: make(chan struct{}),
}
- meta, err := NewMetadata(
+ meta, err := table.NewMetadata(
iceberg.NewSchema(0),
iceberg.UnpartitionedSpec,
- UnsortedSortOrder,
+ table.UnsortedSortOrder,
"s3://bucket/table",
iceberg.Properties{},
)
require.NoError(t, err)
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
"s3://bucket/table/metadata/v1.metadata.json",
testFSF(fsys),
@@ -1891,7 +1895,7 @@ func TestPurgeFilesDeletesNonBulkFilesConcurrently(t
*testing.T) {
)
done := make(chan error, 1)
- go func() { done <- tbl.PurgeFiles(context.Background()) }()
+ go func() { done <- PurgeFiles(context.Background(), tbl) }()
// Every worker remains blocked in Remove until the full pool
is observable.
synctest.Wait()
@@ -1910,12 +1914,12 @@ func
TestPurgeFilesSkipsDataFilesForMalformedGCEnabled(t *testing.T) {
for _, gcValue := range []string{"1", "garbage", "false ", "true "} {
t.Run(gcValue, func(t *testing.T) {
- meta, err := NewMetadata(
+ meta, err := table.NewMetadata(
iceberg.NewSchema(0),
iceberg.UnpartitionedSpec,
- UnsortedSortOrder,
+ table.UnsortedSortOrder,
"s3://bucket/table",
- iceberg.Properties{GCEnabledKey: gcValue},
+ iceberg.Properties{table.GCEnabledKey: gcValue},
)
require.NoError(t, err)
@@ -1923,15 +1927,15 @@ func
TestPurgeFilesSkipsDataFilesForMalformedGCEnabled(t *testing.T) {
path: orphanDataPath,
info: mockFileInfo{name: "orphan.parquet"},
}}}
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
"s3://bucket/table/metadata/v1.metadata.json",
- func(context.Context) (io.IO, error) { return
fsys, nil },
+ func(context.Context) (iceio.IO, error) {
return fsys, nil },
nil,
)
- require.NoError(t, tbl.PurgeFiles(context.Background()))
+ require.NoError(t, PurgeFiles(context.Background(),
tbl))
assert.NotContains(t, fsys.removed, orphanDataPath)
})
}
@@ -1965,20 +1969,20 @@ func TestDeleteOrphanFilesPrefixMismatchModes(t
*testing.T) {
schema := iceberg.NewSchema(0, iceberg.NestedField{
ID: 1, Name: "id", Type:
iceberg.PrimitiveTypes.Int64, Required: true,
})
- meta, err := NewMetadata(
+ meta, err := table.NewMetadata(
schema,
iceberg.UnpartitionedSpec,
- UnsortedSortOrder,
+ table.UnsortedSortOrder,
"s3://bucket/table",
- iceberg.Properties{PropertyFormatVersion: "2"},
+ iceberg.Properties{table.PropertyFormatVersion:
"2"},
)
require.NoError(t, err)
- builder, err := MetadataBuilderFromBase(meta, "")
+ builder, err := table.MetadataBuilderFromBase(meta, "")
require.NoError(t, err)
- require.NoError(t, builder.SetStatistics(StatisticsFile{
+ require.NoError(t,
builder.SetStatistics(table.StatisticsFile{
SnapshotID: 1,
StatisticsPath: tt.referencedPath,
- BlobMetadata: []BlobMetadata{},
+ BlobMetadata: []table.BlobMetadata{},
FileSizeInBytes: 4,
}))
meta, err = builder.Build()
@@ -1990,11 +1994,11 @@ func TestDeleteOrphanFilesPrefixMismatchModes(t
*testing.T) {
info: mockFileInfo{name:
"file.parquet", size: 4},
}},
}
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
"s3://bucket/table/metadata/v1.metadata.json",
- func(context.Context) (io.IO, error) { return
fsys, nil },
+ func(context.Context) (iceio.IO, error) {
return fsys, nil },
nil,
)
@@ -2005,12 +2009,13 @@ func TestDeleteOrphanFilesPrefixMismatchModes(t
*testing.T) {
} {
t.Run(mode.String(), func(t *testing.T) {
var deleted []string
- result, err := tbl.DeleteOrphanFiles(
+ result, err := DeleteOrphanFiles(
context.Background(),
-
WithLocation("s3://bucket/path"),
- WithFilesOlderThan(0),
+ tbl,
+
WithCleanupLocation("s3://bucket/path"),
+ WithCleanupFilesOlderThan(0),
WithPrefixMismatchMode(mode),
- WithDeleteFunc(func(path
string) error {
+ WithCleanupDeleteFunc(func(path
string) error {
deleted =
append(deleted, path)
return nil
@@ -2040,20 +2045,20 @@ func
TestDeleteOrphanFilesDryRunKeepsMixedCaseWindowsReference(t *testing.T) {
schema := iceberg.NewSchema(0, iceberg.NestedField{
ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64,
Required: true,
})
- meta, err := NewMetadata(
+ meta, err := table.NewMetadata(
schema,
iceberg.UnpartitionedSpec,
- UnsortedSortOrder,
+ table.UnsortedSortOrder,
`C:\Warehouse`,
- iceberg.Properties{PropertyFormatVersion: "2"},
+ iceberg.Properties{table.PropertyFormatVersion: "2"},
)
require.NoError(t, err)
- builder, err := MetadataBuilderFromBase(meta, "")
+ builder, err := table.MetadataBuilderFromBase(meta, "")
require.NoError(t, err)
- require.NoError(t, builder.SetStatistics(StatisticsFile{
+ require.NoError(t, builder.SetStatistics(table.StatisticsFile{
SnapshotID: 1,
StatisticsPath: `C:\Warehouse\Data\File.parquet`,
- BlobMetadata: []BlobMetadata{},
+ BlobMetadata: []table.BlobMetadata{},
FileSizeInBytes: 4,
}))
meta, err = builder.Build()
@@ -2066,18 +2071,19 @@ func
TestDeleteOrphanFilesDryRunKeepsMixedCaseWindowsReference(t *testing.T) {
info: mockFileInfo{name: "FILE.PARQUET", size: 4},
}},
}
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
`C:\Warehouse\metadata\v1.metadata.json`,
- func(context.Context) (io.IO, error) { return fsys, nil },
+ func(context.Context) (iceio.IO, error) { return fsys, nil },
nil,
)
- result, err := tbl.DeleteOrphanFiles(
+ result, err := DeleteOrphanFiles(
context.Background(),
- WithFilesOlderThan(0),
- WithDryRun(true),
+ tbl,
+ WithCleanupFilesOlderThan(0),
+ WithCleanupDryRun(true),
)
require.NoError(t, err)
assert.Empty(t, result.OrphanFileLocations)
@@ -2113,20 +2119,20 @@ func
TestDeleteOrphanFilesDryRunKeepsPortableWindowsReferences(t *testing.T) {
schema := iceberg.NewSchema(0, iceberg.NestedField{
ID: 1, Name: "id", Type:
iceberg.PrimitiveTypes.Int64, Required: true,
})
- meta, err := NewMetadata(
+ meta, err := table.NewMetadata(
schema,
iceberg.UnpartitionedSpec,
- UnsortedSortOrder,
+ table.UnsortedSortOrder,
tt.tableLocation,
- iceberg.Properties{PropertyFormatVersion: "2"},
+ iceberg.Properties{table.PropertyFormatVersion:
"2"},
)
require.NoError(t, err)
- builder, err := MetadataBuilderFromBase(meta, "")
+ builder, err := table.MetadataBuilderFromBase(meta, "")
require.NoError(t, err)
- require.NoError(t, builder.SetStatistics(StatisticsFile{
+ require.NoError(t,
builder.SetStatistics(table.StatisticsFile{
SnapshotID: 1,
StatisticsPath: tt.referencedPath,
- BlobMetadata: []BlobMetadata{},
+ BlobMetadata: []table.BlobMetadata{},
FileSizeInBytes: 4,
}))
meta, err = builder.Build()
@@ -2140,18 +2146,19 @@ func
TestDeleteOrphanFilesDryRunKeepsPortableWindowsReferences(t *testing.T) {
})
}
fsys := &mockListableIO{entries: entries}
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
tt.tableLocation+"/metadata/v1.metadata.json",
- func(context.Context) (io.IO, error) { return
fsys, nil },
+ func(context.Context) (iceio.IO, error) {
return fsys, nil },
nil,
)
- result, err := tbl.DeleteOrphanFiles(
+ result, err := DeleteOrphanFiles(
context.Background(),
- WithFilesOlderThan(0),
- WithDryRun(true),
+ tbl,
+ WithCleanupFilesOlderThan(0),
+ WithCleanupDryRun(true),
)
require.NoError(t, err)
assert.Empty(t, result.OrphanFileLocations)
@@ -2165,12 +2172,12 @@ func
TestDeleteOrphanFilesDryRunKeepsOpaqueVersionHint(t *testing.T) {
schema := iceberg.NewSchema(0, iceberg.NestedField{
ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64,
Required: true,
})
- meta, err := NewMetadata(
+ meta, err := table.NewMetadata(
schema,
iceberg.UnpartitionedSpec,
- UnsortedSortOrder,
+ table.UnsortedSortOrder,
tableLocation,
- iceberg.Properties{PropertyFormatVersion: "2"},
+ iceberg.Properties{table.PropertyFormatVersion: "2"},
)
require.NoError(t, err)
@@ -2181,18 +2188,19 @@ func
TestDeleteOrphanFilesDryRunKeepsOpaqueVersionHint(t *testing.T) {
info: mockFileInfo{name: "version-hint.text", size: 2},
}},
}
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
tableLocation+"/metadata/v1.metadata.json",
- func(context.Context) (io.IO, error) { return fsys, nil },
+ func(context.Context) (iceio.IO, error) { return fsys, nil },
nil,
)
- result, err := tbl.DeleteOrphanFiles(
+ result, err := DeleteOrphanFiles(
context.Background(),
- WithFilesOlderThan(0),
- WithDryRun(true),
+ tbl,
+ WithCleanupFilesOlderThan(0),
+ WithCleanupDryRun(true),
)
require.NoError(t, err)
assert.Equal(t, tableLocation, fsys.root)
@@ -2265,7 +2273,7 @@ func TestDeleteOrphanFilesPopulatesOrphanFileSizes(t
*testing.T) {
"refs": {}
}`
- meta, err := ParseMetadataString(metaJSON)
+ meta, err := table.ParseMetadataString(metaJSON)
require.NoError(t, err)
mockFS := &mockListableIO{
@@ -2277,19 +2285,19 @@ func TestDeleteOrphanFilesPopulatesOrphanFileSizes(t
*testing.T) {
},
}
- tbl := New(
- Identifier{"db", "tbl"},
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
"s3://bucket/table/metadata/v1.metadata.json",
- func(context.Context) (io.IO, error) { return mockFS, nil },
+ func(context.Context) (iceio.IO, error) { return mockFS, nil },
nil,
)
- result, err := tbl.DeleteOrphanFiles(context.Background(),
- WithDryRun(true),
- WithLocation("s3://bucket/table"),
+ result, err := DeleteOrphanFiles(context.Background(), tbl,
+ WithCleanupDryRun(true),
+ WithCleanupLocation("s3://bucket/table"),
WithCleanupMaxConcurrency(1),
- WithFilesOlderThan(0), // Consider files created before the
scan.
+ WithCleanupFilesOlderThan(0), // Consider files created before
the scan.
)
require.NoError(t, err)
@@ -2311,20 +2319,26 @@ func TestWalkDirectoryRequiresListableIO(t *testing.T) {
return nil
})
require.Error(t, err)
- assert.Contains(t, err.Error(), "does not implement iceio.ListableIO")
+ assert.Contains(t, err.Error(), "does not implement io.ListableIO")
+}
+
+// testFSF wraps an IO object in a function that returns it.
+// This is a local copy of the helper from table/updates_test.go
+func testFSF(io iceio.IO) func(context.Context) (iceio.IO, error) {
+ return func(context.Context) (iceio.IO, error) { return io, nil }
}
type inMemoryCatalog struct {
- metadata Metadata
+ metadata table.Metadata
}
func (c *inMemoryCatalog) CommitTable(
ctx context.Context,
- ident Identifier,
- reqs []Requirement,
- updates []Update,
-) (Metadata, string, error) {
- meta, err := UpdateTableMetadata(c.metadata, updates, "")
+ ident table.Identifier,
+ reqs []table.Requirement,
+ updates []table.Update,
+) (table.Metadata, string, error) {
+ meta, err := table.UpdateTableMetadata(c.metadata, updates, "")
if err != nil {
return nil, "", err
}
@@ -2333,7 +2347,7 @@ func (c *inMemoryCatalog) CommitTable(
return meta, "", nil
}
-func (c *inMemoryCatalog) LoadTable(ctx context.Context, ident Identifier)
(*Table, error) {
+func (c *inMemoryCatalog) LoadTable(ctx context.Context, ident
table.Identifier) (*table.Table, error) {
return nil, nil
}
@@ -2349,16 +2363,16 @@ func
TestGetReferencedFiles_OverwriteThenExpireExcludesTombstones(t *testing.T)
}, nil)
spec := *iceberg.UnpartitionedSpec
- meta, err := NewMetadata(schema, &spec, UnsortedSortOrder,
tableLocation,
- iceberg.Properties{PropertyFormatVersion: "2"})
+ meta, err := table.NewMetadata(schema, &spec, table.UnsortedSortOrder,
tableLocation,
+ iceberg.Properties{table.PropertyFormatVersion: "2"})
require.NoError(t, err)
- fs := io.LocalFS{}
- tbl := New(
- Identifier{"db", "tbl"},
+ fs := iceio.LocalFS{}
+ tbl := table.New(
+ table.Identifier{"db", "tbl"},
meta,
tableLocation+"/metadata/v0.metadata.json",
- func(context.Context) (io.IO, error) { return fs, nil },
+ func(context.Context) (iceio.IO, error) { return fs, nil },
&inMemoryCatalog{meta},
)
@@ -2394,9 +2408,9 @@ func
TestGetReferencedFiles_OverwriteThenExpireExcludesTombstones(t *testing.T)
// metadata reachability, not the side-effect of file removal.
tx := tbl.NewTransaction()
require.NoError(t, tx.ExpireSnapshots(
- WithRetainLast(1),
- WithOlderThan(0),
- WithPostCommit(false),
+ table.WithRetainLast(1),
+ table.WithOlderThan(0),
+ table.WithPostCommit(false),
))
tbl, err = tx.Commit(ctx)
require.NoError(t, err)
@@ -2405,7 +2419,7 @@ func
TestGetReferencedFiles_OverwriteThenExpireExcludesTombstones(t *testing.T)
// fileA is now referenced only via a DELETED entry in the surviving
// snapshot's tombstone manifest. The fix must exclude it.
- refs, err := tbl.getReferencedFiles(ctx, fs, 1, true)
+ refs, err := getReferencedFiles(ctx, tbl, fs, 1, true)
require.NoError(t, err)
assert.Contains(t, refs, normalizeFilePath(fileB),
@@ -2418,8 +2432,8 @@ func
TestGetReferencedFiles_OverwriteThenExpireExcludesTombstones(t *testing.T)
// given snapshot's manifests, filtered to entries matching wantStatus.
func dataFilePathsFromSnapshot(
t *testing.T,
- snap *Snapshot,
- fs io.IO,
+ snap *table.Snapshot,
+ fs iceio.IO,
wantStatus iceberg.ManifestEntryStatus,
) []string {
t.Helper()
@@ -2456,7 +2470,7 @@ func TestGetReferencedFiles_SharedManifestReadOnce(t
*testing.T) {
writeManifestList(t, tio.trackingIO, 2, manifestList2,
[]iceberg.ManifestFile{mf})
tio.files[dataPath] = []byte("data")
- meta, err := ParseMetadataString(buildMetaJSON(metaJSONOpts{
+ meta, err := table.ParseMetadataString(buildMetaJSON(metaJSONOpts{
snapshots: fmt.Sprintf(
`{"snapshot-id":1,"timestamp-ms":1000,"manifest-list":%q},`+
`{"snapshot-id":2,"timestamp-ms":2000,"manifest-list":%q}`,
@@ -2464,8 +2478,8 @@ func TestGetReferencedFiles_SharedManifestReadOnce(t
*testing.T) {
}))
require.NoError(t, err)
- tbl := New(Identifier{"ns", "tbl"}, meta, "metadata.json",
testFSF(tio), nil)
- refs, err := tbl.getReferencedFiles(context.Background(), tio, 1, true)
+ tbl := table.New(table.Identifier{"ns", "tbl"}, meta, "metadata.json",
testFSF(tio), nil)
+ refs, err := getReferencedFiles(context.Background(), tbl, tio, 1, true)
require.NoError(t, err)
assert.Contains(t, refs, dataPath)
@@ -2495,7 +2509,7 @@ func TestGetReferencedFiles_DisjointManifestsAllRead(t
*testing.T) {
writeManifestList(t, tio.trackingIO, 1, manifestList1,
[]iceberg.ManifestFile{mf1})
writeManifestList(t, tio.trackingIO, 2, manifestList2,
[]iceberg.ManifestFile{mf2})
- meta, err := ParseMetadataString(buildMetaJSON(metaJSONOpts{
+ meta, err := table.ParseMetadataString(buildMetaJSON(metaJSONOpts{
snapshots: fmt.Sprintf(
`{"snapshot-id":1,"timestamp-ms":1000,"manifest-list":%q},`+
`{"snapshot-id":2,"timestamp-ms":2000,"manifest-list":%q}`,
@@ -2503,8 +2517,8 @@ func TestGetReferencedFiles_DisjointManifestsAllRead(t
*testing.T) {
}))
require.NoError(t, err)
- tbl := New(Identifier{"ns", "tbl"}, meta, "metadata.json",
testFSF(tio), nil)
- refs, err := tbl.getReferencedFiles(context.Background(), tio, 1, true)
+ tbl := table.New(table.Identifier{"ns", "tbl"}, meta, "metadata.json",
testFSF(tio), nil)
+ refs, err := getReferencedFiles(context.Background(), tbl, tio, 1, true)
require.NoError(t, err)
assert.Contains(t, refs, dataPath1)
@@ -2534,13 +2548,13 @@ func
TestGetReferencedFiles_ManySnapshotsShareManifest(t *testing.T) {
i, i*1000, listPath))
}
- meta, err := ParseMetadataString(buildMetaJSON(metaJSONOpts{
+ meta, err := table.ParseMetadataString(buildMetaJSON(metaJSONOpts{
snapshots: strings.Join(snapJSON, ","),
}))
require.NoError(t, err)
- tbl := New(Identifier{"ns", "tbl"}, meta, "metadata.json",
testFSF(tio), nil)
- refs, err := tbl.getReferencedFiles(context.Background(), tio, 4, true)
+ tbl := table.New(table.Identifier{"ns", "tbl"}, meta, "metadata.json",
testFSF(tio), nil)
+ refs, err := getReferencedFiles(context.Background(), tbl, tio, 4, true)
require.NoError(t, err)
assert.Contains(t, refs, dataPath)
@@ -2568,14 +2582,14 @@ func
TestGetReferencedFilesFetchesManifestListsConcurrently(t *testing.T) {
i, i*1000, listPath))
}
- meta, err := ParseMetadataString(buildMetaJSON(metaJSONOpts{
+ meta, err := table.ParseMetadataString(buildMetaJSON(metaJSONOpts{
snapshots: strings.Join(snapshotJSON, ","),
}))
require.NoError(t, err)
trackingFS := &manifestTrackingIO{IO: baseIO, delay: 10 *
time.Millisecond}
- tbl := New(Identifier{"ns", "tbl"}, meta, "metadata.json",
testFSF(trackingFS), nil)
- refs, err := tbl.getReferencedFiles(context.Background(), trackingFS,
maxWorkers, true)
+ tbl := table.New(table.Identifier{"ns", "tbl"}, meta, "metadata.json",
testFSF(trackingFS), nil)
+ refs, err := getReferencedFiles(context.Background(), tbl, trackingFS,
maxWorkers, true)
require.NoError(t, err)
assert.Contains(t, refs, dataPath)
diff --git a/table/maintenance/retention_validation_test.go
b/table/maintenance/retention_validation_test.go
new file mode 100644
index 000000000..003157a67
--- /dev/null
+++ b/table/maintenance/retention_validation_test.go
@@ -0,0 +1,38 @@
+// 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 maintenance
+
+import (
+ "context"
+ "testing"
+ "time"
+
+ "github.com/apache/iceberg-go/table"
+ "github.com/stretchr/testify/require"
+)
+
+func TestDeleteOrphanFilesRejectsNegativeAgeBeforeScan(t *testing.T) {
+ _, err := DeleteOrphanFiles(context.Background(), &table.Table{},
WithCleanupFilesOlderThan(-time.Nanosecond))
+ require.EqualError(t, err, "orphan cleanup age must be non-negative")
+}
+
+func TestRetentionOptionsAcceptZeroAge(t *testing.T) {
+ cfg := &orphanCleanupConfig{}
+ WithCleanupFilesOlderThan(0)(cfg)
+ require.NoError(t, cfg.validationErr)
+}
diff --git a/table/properties.go b/table/properties.go
index 03573bee2..0743e660f 100644
--- a/table/properties.go
+++ b/table/properties.go
@@ -176,7 +176,15 @@ const (
CommitTotalRetryTimeoutMsDefault = 30 * 60 * 1000
)
-func isGCEnabled(props iceberg.Properties) bool {
+// IsGCEnabled reports whether garbage collection (physical file deletion) is
+// enabled for the given table properties. A missing "gc.enabled" key returns
+// [GCEnabledDefault]. Otherwise GC is enabled only when the value is a
+// case-insensitive "true"; any other value, including "1", "true " (with
+// surrounding whitespace), or arbitrary text, disables it.
+//
+// This is exported so the table/maintenance package can gate destructive
+// operations on it; the parsing rules above are part of that contract.
+func IsGCEnabled(props iceberg.Properties) bool {
value, ok := props[GCEnabledKey]
if !ok {
return GCEnabledDefault
diff --git a/table/retention_validation_test.go
b/table/retention_validation_test.go
index dd3b2817c..8b208b561 100644
--- a/table/retention_validation_test.go
+++ b/table/retention_validation_test.go
@@ -18,18 +18,12 @@
package table
import (
- "context"
"testing"
"time"
"github.com/stretchr/testify/require"
)
-func TestDeleteOrphanFilesRejectsNegativeAgeBeforeScan(t *testing.T) {
- _, err := (Table{}).DeleteOrphanFiles(context.Background(),
WithFilesOlderThan(-time.Nanosecond))
- require.EqualError(t, err, "orphan cleanup age must be non-negative")
-}
-
func TestExpireSnapshotsRejectsInvalidRetentionBeforeMetadataAccess(t
*testing.T) {
tests := []struct {
name string
@@ -51,10 +45,6 @@ func
TestExpireSnapshotsRejectsInvalidRetentionBeforeMetadataAccess(t *testing.T
}
func TestRetentionOptionsAcceptZeroAgeAndOneSnapshot(t *testing.T) {
- cleanupCfg := &orphanCleanupConfig{}
- WithFilesOlderThan(0)(cleanupCfg)
- require.NoError(t, cleanupCfg.validationErr)
-
expireCfg := &expireSnapshotsCfg{}
WithOlderThan(0)(expireCfg)
WithRetainLast(1)(expireCfg)
diff --git a/table/transaction.go b/table/transaction.go
index 94817657b..18d9d5d89 100644
--- a/table/transaction.go
+++ b/table/transaction.go
@@ -3393,7 +3393,7 @@ func (t *Transaction) Commit(ctx context.Context)
(*Table, error) {
}
func validateGCEnabledForSnapshotExpiration(props iceberg.Properties) error {
- if !isGCEnabled(props) {
+ if !IsGCEnabled(props) {
return errors.New("cannot expire snapshots: GC is disabled
(deleting files may corrupt other tables)")
}
diff --git a/table/transaction_internal_test.go
b/table/transaction_internal_test.go
index b9eedc2db..306916335 100644
--- a/table/transaction_internal_test.go
+++ b/table/transaction_internal_test.go
@@ -912,7 +912,7 @@ func TestIsGCEnabled(t *testing.T) {
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
- require.Equal(t, tt.want, isGCEnabled(tt.props))
+ require.Equal(t, tt.want, IsGCEnabled(tt.props))
})
}
}
diff --git a/table/updates.go b/table/updates.go
index b9ad36b0c..fa50ac034 100644
--- a/table/updates.go
+++ b/table/updates.go
@@ -657,7 +657,7 @@ func (u *removeSnapshotsUpdate) Apply(builder
*MetadataBuilder) error {
}
func (u *removeSnapshotsUpdate) PostCommit(ctx context.Context, preTable
*Table, postTable *Table) error {
- if !u.postCommit || postTable == nil ||
!isGCEnabled(postTable.Properties()) {
+ if !u.postCommit || postTable == nil ||
!IsGCEnabled(postTable.Properties()) {
return nil
}
diff --git a/website/src/api.md b/website/src/api.md
index 83747beb6..225c06b01 100644
--- a/website/src/api.md
+++ b/website/src/api.md
@@ -484,19 +484,19 @@ Tune retention with the
`history.expire.min-snapshots-to-keep`, `history.expire.
import (
"time"
- "github.com/apache/iceberg-go/table"
+ "github.com/apache/iceberg-go/table/maintenance"
)
-result, err := tbl.DeleteOrphanFiles(ctx,
- table.WithFilesOlderThan(72*time.Hour),
- table.WithDryRun(false),
- table.WithCleanupMaxConcurrency(8),
+result, err := maintenance.DeleteOrphanFiles(ctx, tbl,
+ maintenance.WithCleanupFilesOlderThan(72*time.Hour),
+ maintenance.WithCleanupDryRun(false),
+ maintenance.WithCleanupMaxConcurrency(8),
)
if err != nil { /* ... */ }
fmt.Printf("removed %d files\n", len(result.DeletedFiles))
```
-Also see `table.WithLocation`, `table.WithDeleteFunc`,
`table.WithPrefixMismatchMode`, `table.WithEqualSchemes`, and
`table.WithEqualAuthorities` in `table/orphan_cleanup.go`.
+Also see `maintenance.WithCleanupLocation`,
`maintenance.WithCleanupDeleteFunc`, `maintenance.WithPrefixMismatchMode`,
`maintenance.WithEqualSchemes`, and `maintenance.WithEqualAuthorities` in
`table/maintenance/orphan_cleanup.go`.
### Compaction (rewrite data files)
diff --git a/website/src/releases.md b/website/src/releases.md
index 59e548a51..eab3fa540 100644
--- a/website/src/releases.md
+++ b/website/src/releases.md
@@ -42,7 +42,14 @@ Per-release notes (highlights, breaking changes,
contributors) are published on
### Unreleased compatibility notes
-- `Table.PurgeFiles` now calls `IO.Remove` concurrently, with up to 32 calls
in flight, when the filesystem does not implement `BulkRemovableIO`. Custom IO
implementations, including those registered with `io.Register`, must make
`Remove` safe for concurrent use and synchronize any shared mutable state. The
bulk deletion path is unchanged.
+- **Orphan file cleanup moved to `table/maintenance`.** The orphan-cleanup and
file-purge API is no longer on `table.Table`; it now lives in the
`table/maintenance` package as free functions that take the table as their
first argument (matching `table/compaction`). For example,
`tbl.DeleteOrphanFiles(ctx, ...)` becomes `maintenance.DeleteOrphanFiles(ctx,
tbl, ...)`. Because `table/maintenance` imports `table` (and not the reverse),
this cannot be offered as a deprecation shim on `table.T [...]
+ - Functions (now free functions taking `*table.Table`):
`DeleteOrphanFiles`, `PlanOrphanFiles`, `ExecuteOrphanCleanup`, `PurgeFiles`.
+ - Types: `OrphanCleanupOption`, `OrphanCleanupResult`, `OrphanFile`,
`OrphanCleanupPlan`, `PrefixMismatchMode`.
+ - Constants: `PrefixMismatchError`, `PrefixMismatchIgnore`,
`PrefixMismatchDelete`.
+ - Option constructors: `WithCleanupLocation`, `WithCleanupFilesOlderThan`,
`WithCleanupDryRun`, `WithCleanupDeleteFunc`, `WithCleanupMaxConcurrency`,
`WithPrefixMismatchMode`, `WithEqualSchemes`, `WithEqualAuthorities`. The four
generic names (`WithLocation`, `WithFilesOlderThan`, `WithDryRun`,
`WithDeleteFunc`) gained the `WithCleanup` prefix for consistency.
+
+ Relatedly, the previously unexported `isGCEnabled` is now
`table.IsGCEnabled`, so the new package can gate destructive operations on the
`gc.enabled` table property.
+- `maintenance.PurgeFiles` (previously `Table.PurgeFiles`) now calls
`IO.Remove` concurrently, with up to 32 calls in flight, when the filesystem
does not implement `BulkRemovableIO`. Custom IO implementations, including
those registered with `io.Register`, must make `Remove` safe for concurrent use
and synchronize any shared mutable state. The bulk deletion path is unchanged.
## Verifying and producing releases