zeroshade commented on code in PR #2086:
URL: https://github.com/apache/iceberg-go/pull/2086#discussion_r4160434456
##########
table/maintenance/orphan_cleanup_test.go:
##########
@@ -2438,152 +2381,3 @@ func dataFilePathsFromSnapshot(
return paths
}
Review Comment:
Four tests were deleted here with no replacement:
`TestGetReferencedFiles_SharedManifestReadOnce`, `_DisjointManifestsAllRead`,
`_ManySnapshotsShareManifest` and
`TestGetReferencedFilesFetchesManifestListsConcurrently` (old
`table/orphan_cleanup_test.go:2442-2589`). They were the only tests checking
that a manifest shared across snapshots is opened once, and that manifest-list
fetches stay within the concurrency limit (`1 < maxOpen <= maxWorkers`).
The helpers they need are test-only in `table`:
`trackingCallsIO`/`newTrackingCallsIO`, `writeManifest`, `writeManifestList`
and `metaJSONOpts`/`buildMetaJSON` (`table/updates_test.go:39-140`),
`trackingIO`/`newTrackingIO` (`table/snapshot_producers_test.go:698-710`), and
`manifestTrackingIO` (`table/all_manifests_internal_test.go:129`). Copy them
here like you did with `testFSF`, or move them into a shared internal
test-helper package, and restore the four tests.
##########
table/maintenance/orphan_cleanup_test.go:
##########
@@ -1444,69 +1443,77 @@ func TestOrphanCleanup_EdgeCases(t *testing.T) {
})
}
-func TestGetReferencedFiles_IncludesStatisticsFiles(t *testing.T) {
- const metaJSON = `{
- "format-version": 2,
- "table-uuid": "9c12d441-03fe-4693-9a96-a0705ddf69c1",
- "location": "s3://bucket/test/location",
- "last-sequence-number": 0,
- "last-updated-ms": 1602638573590,
- "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": []}],
- "metadata-log": [],
- "snapshot-log": [],
- "statistics": [
- {
- "snapshot-id": 1,
- "statistics-path": "s3://bucket/stats/table-stats.puffin",
- "file-size-in-bytes": 1024,
- "file-footer-size-in-bytes": 512,
- "blob-metadata": []
- },
- {
- "snapshot-id": 2,
- "statistics-path": "",
- "file-size-in-bytes": 0,
- "file-footer-size-in-bytes": 0,
- "blob-metadata": []
- }
- ],
- "partition-statistics": [
- {
- "snapshot-id": 1,
- "statistics-path": "s3://bucket/stats/part-stats.puffin",
- "file-size-in-bytes": 512
- }
- ]
-}`
-
- meta, err := ParseMetadataString(metaJSON)
- require.NoError(t, err)
+// This test is skipped as it tests private implementation details of
getReferencedFiles,
+// which is now part of the Service internal API. The public API
(PlanOrphanFiles, PurgeFiles)
+// is tested through other tests.
+func SkipTestGetReferencedFiles_IncludesStatisticsFiles(t *testing.T) {
+ // This test is skipped because it tests private implementation details
+ // of getReferencedFiles, which is now an internal method on Service.
+ t.Skip("tests private implementation details of getReferencedFiles")
+ /*
Review Comment:
`go test` only discovers `Test*` functions, so `SkipTest*` never runs and
doesn't even report SKIP. On this head, `go test -run GetReferencedFiles
./table/maintenance` prints "no tests to run". The stated reason doesn't hold
either: this file is `package maintenance` and already calls unexported helpers
(`newOrphanCleanupConfig`, `walkDirectory`, `deleteFilesParallel`), so
`New(tbl).getReferencedFiles(ctx, nil, 1, true)` is callable from here.
Please port it: build the table with `table.New(table.Identifier{"db",
"tbl"}, meta, "s3://bucket/test/location/metadata/v1.metadata.json", nil,
nil)`, use `tbl.MetadataLocation()` instead of `tbl.metadataLocation`, and
delete the commented-out block. As it stands, nothing checks that
partition-statistics paths count as referenced, and that check is what stops
`DeleteOrphanFiles` from deleting live partition-stats files.
##########
table/maintenance/orphan_cleanup_bench_test.go:
##########
@@ -0,0 +1,65 @@
+// 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 (
+ "fmt"
+ "testing"
+)
+
+var orphanCleanupBenchmarkSink string
+
+// BenchmarkApplyURIEquivalence measures the per-path lookup cost as the
+// number of paths and configured equivalence groups grows. Configuration is
+// built before the timer starts so only the repeated lookup is measured.
+func BenchmarkApplyURIEquivalence(b *testing.B) {
+ for _, tc := range []struct {
+ paths int
+ groups int
+ }{
+ {paths: 100, groups: 1},
+ {paths: 100, groups: 100},
+ {paths: 10_000, groups: 1},
+ {paths: 10_000, groups: 100},
+ } {
+ b.Run(fmt.Sprintf("paths=%d/groups=%d", tc.paths, tc.groups),
func(b *testing.B) {
+ equivalences := make(map[string]string, tc.groups)
+ for i := range tc.groups {
+
equivalences[fmt.Sprintf("scheme-%d,scheme-%d-alt", i, i)] = "canonical"
+ }
+ cfg :=
newOrphanCleanupConfig(WithEqualSchemes(equivalences))
+
+ schemes := make([]string, tc.paths)
+ for i := range schemes {
+ schemes[i] = fmt.Sprintf("scheme-%d",
i%tc.groups)
+ }
+
+ b.ReportAllocs()
+ var result string
+ b.ResetTimer()
+ for b.Loop() {
+ for _, scheme := range schemes {
+ result = applySchemeEquivalence(scheme,
cfg.equalSchemes)
+ }
+ }
+ b.StopTimer()
+ // Keep the result observable without including the
sink in the benchmark.
+ orphanCleanupBenchmarkSink = result
+ })
+ }
+}
Review Comment:
`BenchmarkPurgeFilesNonBulkDeletion` and
`BenchmarkGetReferencedFilesManifestLists` (with `benchmarkDelayIO`) were
dropped in the move. The PurgeFiles one only calls `deleteFilesParallel`, which
now lives in this package, so it can move unchanged; it's the measurement
behind the 32-worker purge default from #1973. The manifest-list one needs the
same helpers as the deleted dedup tests.
##########
table/maintenance/maintenance.go:
##########
@@ -0,0 +1,30 @@
+// 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 "github.com/apache/iceberg-go/table"
+
+// Service provides maintenance operations on a table, including orphan file
cleanup.
+type Service struct {
+ tbl *table.Table
+}
+
+// New creates a new maintenance service for the given table.
+func New(tbl *table.Table) *Service {
+ return &Service{tbl: tbl}
+}
Review Comment:
#2083 lists the shape (free functions vs facade) as an open question to
settle before slice 1 lands, and nobody has discussed it on the issue yet.
`table/compaction` already uses free functions (`compaction.Analyze(ctx, tbl,
cfg)`), and `Service` here wraps one field and holds no state. If
`New(*table.Table) *Service` is frozen at v1, adding facade-level config later
(metrics reporter, executor) is a breaking change. Let's choose on #2083
between free functions and `New(tbl, opts ...Option)`, and use the same shape
for the inspect slice.
The package also needs a `// Package maintenance` doc comment
(`compaction.go:18` has one). And since this package is meant to grow more
actions, consider whether `WithLocation`, `WithDryRun`, `WithDeleteFunc` and
`WithFilesOlderThan` should take the `WithCleanup*` prefix that
`WithCleanupMaxConcurrency` already has, before v1 freezes the names.
##########
table/maintenance/orphan_cleanup_test.go:
##########
@@ -2333,93 +2346,23 @@ 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
}
-func TestGetReferencedFiles_OverwriteThenExpireExcludesTombstones(t
*testing.T) {
- ctx := context.Background()
- tableLocation := t.TempDir()
-
- schema := iceberg.NewSchema(0,
- iceberg.NestedField{ID: 1, Name: "id", Type:
iceberg.PrimitiveTypes.Int64, Required: true},
- )
- arrSchema := arrow.NewSchema([]arrow.Field{
- {Name: "id", Type: arrow.PrimitiveTypes.Int64, Nullable: false},
- }, nil)
- spec := *iceberg.UnpartitionedSpec
-
- meta, err := NewMetadata(schema, &spec, UnsortedSortOrder,
tableLocation,
- iceberg.Properties{PropertyFormatVersion: "2"})
- require.NoError(t, err)
-
- fs := io.LocalFS{}
- tbl := New(
- Identifier{"db", "tbl"},
- meta,
- tableLocation+"/metadata/v0.metadata.json",
- func(context.Context) (io.IO, error) { return fs, nil },
- &inMemoryCatalog{meta},
- )
-
- // Step 1: append id=1. Produces snapshot 1 with one ADDED entry for
fileA.
- arrA, err := array.TableFromJSON(memory.DefaultAllocator, arrSchema,
[]string{`[{"id": 1}]`})
- require.NoError(t, err)
- defer arrA.Release()
- tbl, err = tbl.AppendTable(ctx, arrA, 1, nil)
- require.NoError(t, err)
-
- snap1 := tbl.CurrentSnapshot()
- require.NotNil(t, snap1)
- pathA := dataFilePathsFromSnapshot(t, snap1, fs,
iceberg.EntryStatusADDED)
- require.Len(t, pathA, 1, "expected one ADDED data file after append")
- fileA := pathA[0]
-
- // Step 2: overwrite with id=2. Produces snapshot 2 whose manifest list
- // contains [added-fileB-manifest, deleted-fileA-manifest]. fileA still
- // lives in snapshot 1's manifest as ADDED at this point.
- arrB, err := array.TableFromJSON(memory.DefaultAllocator, arrSchema,
[]string{`[{"id": 2}]`})
- require.NoError(t, err)
- defer arrB.Release()
- tbl, err = tbl.OverwriteTable(ctx, arrB, 1, nil)
- require.NoError(t, err)
- require.Len(t, tbl.Metadata().Snapshots(), 2, "expected two snapshots
after overwrite")
-
- pathB := dataFilePathsFromSnapshot(t, tbl.CurrentSnapshot(), fs,
iceberg.EntryStatusADDED)
- require.Len(t, pathB, 1, "expected one ADDED data file after overwrite")
- fileB := pathB[0]
-
- // Step 3: expire snapshot 1, keeping only the overwrite snapshot.
- // WithPostCommit(false) keeps fileA on disk so the test only exercises
- // metadata reachability, not the side-effect of file removal.
- tx := tbl.NewTransaction()
- require.NoError(t, tx.ExpireSnapshots(
- WithRetainLast(1),
- WithOlderThan(0),
- WithPostCommit(false),
- ))
- tbl, err = tx.Commit(ctx)
- require.NoError(t, err)
- require.Len(t, tbl.Metadata().Snapshots(), 1,
- "only the overwrite snapshot should remain after expiration")
-
- // 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)
- require.NoError(t, err)
-
- assert.Contains(t, refs, normalizeFilePath(fileB),
- "new live file (ADDED in surviving snapshot) must be in
reference set")
- assert.NotContains(t, refs, normalizeFilePath(fileA),
- "overwritten file (only present as DELETED tombstone) must NOT
be in reference set")
+func SkipTestGetReferencedFiles_OverwriteThenExpireExcludesTombstones(t
*testing.T) {
+ // This test is skipped because it tests private implementation details
+ // of getReferencedFiles, which is now an internal method on Service.
+ // The public API (PlanOrphanFiles, PurgeFiles) is tested through other
tests.
+ t.Skip("tests private implementation details of getReferencedFiles")
}
-// dataFilePathsFromSnapshot returns the data-file paths referenced by the
+// dataFilePathsFromtable.Snapshot returns the data-file paths referenced by
the
Review Comment:
The find-and-replace garbled this comment
(`dataFilePathsFromtable.Snapshot`). Also, `dataFilePathsFromSnapshot` and
`inMemoryCatalog` (L2330) have no callers now that the tombstone test is a
stub. The linter doesn't flag it because `unused` isn't enabled in
`.golangci.yml`. Restoring the tombstone test fixes both; otherwise delete them.
##########
table/maintenance/orphan_cleanup_test.go:
##########
@@ -2333,93 +2346,23 @@ 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
}
-func TestGetReferencedFiles_OverwriteThenExpireExcludesTombstones(t
*testing.T) {
- ctx := context.Background()
- tableLocation := t.TempDir()
-
- schema := iceberg.NewSchema(0,
- iceberg.NestedField{ID: 1, Name: "id", Type:
iceberg.PrimitiveTypes.Int64, Required: true},
- )
- arrSchema := arrow.NewSchema([]arrow.Field{
- {Name: "id", Type: arrow.PrimitiveTypes.Int64, Nullable: false},
- }, nil)
- spec := *iceberg.UnpartitionedSpec
-
- meta, err := NewMetadata(schema, &spec, UnsortedSortOrder,
tableLocation,
- iceberg.Properties{PropertyFormatVersion: "2"})
- require.NoError(t, err)
-
- fs := io.LocalFS{}
- tbl := New(
- Identifier{"db", "tbl"},
- meta,
- tableLocation+"/metadata/v0.metadata.json",
- func(context.Context) (io.IO, error) { return fs, nil },
- &inMemoryCatalog{meta},
- )
-
- // Step 1: append id=1. Produces snapshot 1 with one ADDED entry for
fileA.
- arrA, err := array.TableFromJSON(memory.DefaultAllocator, arrSchema,
[]string{`[{"id": 1}]`})
- require.NoError(t, err)
- defer arrA.Release()
- tbl, err = tbl.AppendTable(ctx, arrA, 1, nil)
- require.NoError(t, err)
-
- snap1 := tbl.CurrentSnapshot()
- require.NotNil(t, snap1)
- pathA := dataFilePathsFromSnapshot(t, snap1, fs,
iceberg.EntryStatusADDED)
- require.Len(t, pathA, 1, "expected one ADDED data file after append")
- fileA := pathA[0]
-
- // Step 2: overwrite with id=2. Produces snapshot 2 whose manifest list
- // contains [added-fileB-manifest, deleted-fileA-manifest]. fileA still
- // lives in snapshot 1's manifest as ADDED at this point.
- arrB, err := array.TableFromJSON(memory.DefaultAllocator, arrSchema,
[]string{`[{"id": 2}]`})
- require.NoError(t, err)
- defer arrB.Release()
- tbl, err = tbl.OverwriteTable(ctx, arrB, 1, nil)
- require.NoError(t, err)
- require.Len(t, tbl.Metadata().Snapshots(), 2, "expected two snapshots
after overwrite")
-
- pathB := dataFilePathsFromSnapshot(t, tbl.CurrentSnapshot(), fs,
iceberg.EntryStatusADDED)
- require.Len(t, pathB, 1, "expected one ADDED data file after overwrite")
- fileB := pathB[0]
-
- // Step 3: expire snapshot 1, keeping only the overwrite snapshot.
- // WithPostCommit(false) keeps fileA on disk so the test only exercises
- // metadata reachability, not the side-effect of file removal.
- tx := tbl.NewTransaction()
- require.NoError(t, tx.ExpireSnapshots(
- WithRetainLast(1),
- WithOlderThan(0),
- WithPostCommit(false),
- ))
- tbl, err = tx.Commit(ctx)
- require.NoError(t, err)
- require.Len(t, tbl.Metadata().Snapshots(), 1,
- "only the overwrite snapshot should remain after expiration")
-
- // 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)
- require.NoError(t, err)
-
- assert.Contains(t, refs, normalizeFilePath(fileB),
- "new live file (ADDED in surviving snapshot) must be in
reference set")
- assert.NotContains(t, refs, normalizeFilePath(fileA),
- "overwritten file (only present as DELETED tombstone) must NOT
be in reference set")
+func SkipTestGetReferencedFiles_OverwriteThenExpireExcludesTombstones(t
*testing.T) {
+ // This test is skipped because it tests private implementation details
+ // of getReferencedFiles, which is now an internal method on Service.
+ // The public API (PlanOrphanFiles, PurgeFiles) is tested through other
tests.
+ t.Skip("tests private implementation details of getReferencedFiles")
}
Review Comment:
Same problem: this never runs, and it was the only test for excluding
DELETED tombstones in `getReferencedFiles`. Everything it needs is exported or
already ported in this file: `tbl.AppendTable`, `tbl.OverwriteTable`,
`tbl.NewTransaction().ExpireSnapshots(table.WithRetainLast(1),
table.WithOlderThan(0), table.WithPostCommit(false))`, plus `inMemoryCatalog`
and `dataFilePathsFromSnapshot`. Please restore the original body, calling
`New(tbl).getReferencedFiles(ctx, fs, 1, true)`.
##########
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) {
Review Comment:
`TestDeleteOrphanFilesRejectsNegativeAgeBeforeScan` and the
`WithFilesOlderThan(0)` half of
`TestRetentionOptionsAcceptZeroAgeAndOneSnapshot` were removed and never
re-added in `table/maintenance`. Nothing tests `orphan cleanup age must be
non-negative` (`maintenance/orphan_cleanup.go:111`) anymore. In the maintenance
tests it is one line:
```go
_, err := New(&table.Table{}).DeleteOrphanFiles(context.Background(),
WithFilesOlderThan(-time.Nanosecond))
require.EqualError(t, err, "orphan cleanup age must be non-negative")
```
##########
table/properties.go:
##########
@@ -176,7 +176,8 @@ const (
CommitTotalRetryTimeoutMsDefault = 30 * 60 * 1000
)
-func isGCEnabled(props iceberg.Properties) bool {
+// IsGCEnabled checks if garbage collection is enabled for the table
properties.
+func IsGCEnabled(props iceberg.Properties) bool {
Review Comment:
This adds public API to `table` during the freeze window only so
`maintenance` can use it. If it stays exported, the doc comment should state
the parsing rules, since they become a contract: a missing key means
`GCEnabledDefault`, and only a case-insensitive `"true"` enables GC (`"1"`,
`"true "` and `"garbage"` all disable it;
`TestPurgeFilesSkipsDataFilesForMalformedGCEnabled` relies on this). Otherwise
move it and the key into `table/internal`, which `table/maintenance` can import.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]