diff --git a/weed/plugin/worker/iceberg/detection.go b/weed/plugin/worker/iceberg/detection.go index 4523af10c..d97ce62f1 100644 --- a/weed/plugin/worker/iceberg/detection.go +++ b/weed/plugin/worker/iceberg/detection.go @@ -333,20 +333,15 @@ func hasEligibleCompaction( } var dataManifests []iceberg.ManifestFile - specIDs := make(map[int32]struct{}) for _, mf := range manifests { if mf.ManifestContent() != iceberg.ManifestContentData { continue } dataManifests = append(dataManifests, mf) - specIDs[mf.PartitionSpecID()] = struct{}{} } if len(dataManifests) == 0 { return false, nil } - if len(specIDs) > 1 { - return false, nil - } var allEntries []iceberg.ManifestEntry for _, mf := range dataManifests { diff --git a/weed/plugin/worker/iceberg/exec_test.go b/weed/plugin/worker/iceberg/exec_test.go index cbb4c088a..196b0d05b 100644 --- a/weed/plugin/worker/iceberg/exec_test.go +++ b/weed/plugin/worker/iceberg/exec_test.go @@ -410,6 +410,19 @@ func populateTable(t *testing.T, fs *fakeFilerServer, setup tableSetup) table.Me func writeCurrentSnapshotManifests(t *testing.T, fs *fakeFilerServer, setup tableSetup, meta table.Metadata, manifestEntries [][]iceberg.ManifestEntry) { t.Helper() + specs := make([]iceberg.PartitionSpec, len(manifestEntries)) + for i := range specs { + specs[i] = meta.PartitionSpec() + } + writeCurrentSnapshotManifestsWithSpecs(t, fs, setup, meta, specs, manifestEntries) +} + +func writeCurrentSnapshotManifestsWithSpecs(t *testing.T, fs *fakeFilerServer, setup tableSetup, meta table.Metadata, specs []iceberg.PartitionSpec, manifestEntries [][]iceberg.ManifestEntry) { + t.Helper() + if len(specs) != len(manifestEntries) { + t.Fatalf("spec count %d does not match manifest count %d", len(specs), len(manifestEntries)) + } + currentSnap := meta.CurrentSnapshot() if currentSnap == nil { t.Fatal("current snapshot is required") @@ -418,7 +431,6 @@ func writeCurrentSnapshotManifests(t *testing.T, fs *fakeFilerServer, setup tabl metaDir := path.Join(s3tables.TablesPath, setup.BucketName, setup.tablePath(), "metadata") version := meta.Version() schema := meta.CurrentSchema() - spec := meta.PartitionSpec() var manifests []iceberg.ManifestFile for i, entries := range manifestEntries { @@ -428,7 +440,7 @@ func writeCurrentSnapshotManifests(t *testing.T, fs *fakeFilerServer, setup tabl path.Join("metadata", manifestName), &manifestBuf, version, - spec, + specs[i], schema, currentSnap.SnapshotID, entries, @@ -1175,6 +1187,58 @@ func TestDetectSchedulesCompactionWithDeleteManifestPresent(t *testing.T) { } } +func TestDetectSchedulesCompactionForSpecEvolvedTable(t *testing.T) { + fs, client := startFakeFiler(t) + + now := time.Now().UnixMilli() + setup := tableSetup{ + BucketName: "test-bucket", + Namespace: "analytics", + TableName: "events", + Snapshots: []table.Snapshot{ + {SnapshotID: 1, TimestampMs: now, ManifestList: "metadata/snap-1.avro", SequenceNumber: 1}, + }, + } + meta := populateTable(t, fs, setup) + + spec0 := meta.PartitionSpec() + spec1 := iceberg.NewPartitionSpecID(1) + partSpecs := map[int32]iceberg.PartitionSpec{ + 0: spec0, + 1: spec1, + } + writeCurrentSnapshotManifestsWithSpecs(t, fs, setup, meta, + []iceberg.PartitionSpec{spec0, spec1}, + [][]iceberg.ManifestEntry{ + makeTestEntriesWithSpec(t, []testEntrySpec{ + {path: "data/spec0-small-1.parquet", size: 1024, partition: map[int]any{}, specID: 0}, + {path: "data/spec0-small-2.parquet", size: 1024, partition: map[int]any{}, specID: 0}, + }, partSpecs), + makeTestEntriesWithSpec(t, []testEntrySpec{ + {path: "data/spec1-small-1.parquet", size: 1024, partition: map[int]any{}, specID: 1}, + {path: "data/spec1-small-2.parquet", size: 1024, partition: map[int]any{}, specID: 1}, + }, partSpecs), + }, + ) + + handler := NewHandler(nil) + config := Config{ + SnapshotRetentionHours: 24 * 365, + MaxSnapshotsToKeep: 10, + TargetFileSizeBytes: 4096, + MinInputFiles: 2, + Operations: "compact", + } + + tables, err := handler.scanTablesForMaintenance(context.Background(), client, config, "", "", "", 0) + if err != nil { + t.Fatalf("scanTablesForMaintenance failed: %v", err) + } + if len(tables) != 1 { + t.Fatalf("expected 1 compaction candidate for spec-evolved table, got %d", len(tables)) + } +} + func TestDetectSchedulesSnapshotExpiryDespiteCompactionEvaluationError(t *testing.T) { fs, client := startFakeFiler(t)