From 53df5428ca43201787b5f3896cc29110a0317e66 Mon Sep 17 00:00:00 2001 From: badalprasadsingh Date: Fri, 25 Sep 2026 17:33:07 +0530 Subject: [PATCH 1/2] fix Signed-off-by: badalprasadsingh --- table/dv_rewrite_test.go | 10 +- table/replace_files_test.go | 436 +++++++++++++++++++++++++++-- table/rewrite_files_test.go | 6 +- table/row_delta.go | 5 +- table/snapshots.go | 7 +- table/transaction.go | 57 ++-- table/transaction_internal_test.go | 5 +- 7 files changed, 458 insertions(+), 68 deletions(-) diff --git a/table/dv_rewrite_test.go b/table/dv_rewrite_test.go index 5620b87f1..abbeed898 100644 --- a/table/dv_rewrite_test.go +++ b/table/dv_rewrite_test.go @@ -225,7 +225,7 @@ func TestRewriteFiles_ExplicitDeleteRewriteIgnoresAutomaticDVRemoval(t *testing. tbl, err = tx.Commit(ctx) require.NoError(t, err) - assert.Equal(t, oldEqualitySequence, deleteFileSequence(t, tbl, newEqualityPath)) + assert.Equal(t, oldEqualitySequence, currentManifestEntry(t, tbl, newEqualityPath).SequenceNum()) assert.Empty(t, deleteEntriesReferencing(t, tbl, rewritten)) assertRowCount(t, tbl, 3) } @@ -246,7 +246,7 @@ func TestRewriteFiles_DataSequenceNumberKeepsDVApplicableToReplacement(t *testin require.NoError(t, err) require.Len(t, tasks, 1) oldData := tasks[0].File - oldDataSequence := fileDataSequence(t, tbl, oldDataPath) + oldDataSequence := currentManifestEntry(t, tbl, oldDataPath).SequenceNum() oldDeletePath := tbl.Location() + "/data/old-pos-delete.parquet" writeParquetFile(t, oldDeletePath, table.PositionalDeleteArrowSchema, @@ -260,7 +260,7 @@ func TestRewriteFiles_DataSequenceNumberKeepsDVApplicableToReplacement(t *testin require.NoError(t, tx.NewRowDelta(nil).AddDeletes(oldDelete).Commit(t.Context())) tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - oldDeleteSequence := fileDataSequence(t, tbl, oldDeletePath) + oldDeleteSequence := currentManifestEntry(t, tbl, oldDeletePath).SequenceNum() tx = tbl.NewTransaction() require.NoError(t, tx.UpgradeFormatVersion(3)) @@ -302,8 +302,8 @@ func TestRewriteFiles_DataSequenceNumberKeepsDVApplicableToReplacement(t *testin tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - assert.Equal(t, oldDataSequence, fileDataSequence(t, tbl, newDataPath)) - assert.Equal(t, oldDeleteSequence, fileDataSequence(t, tbl, newDVs[0].FilePath())) + assert.Equal(t, oldDataSequence, currentManifestEntry(t, tbl, newDataPath).SequenceNum()) + assert.Equal(t, oldDeleteSequence, currentManifestEntry(t, tbl, newDVs[0].FilePath()).SequenceNum()) assert.Equal(t, []int64{2}, scanIDs(t, tbl), "the deletion vector must remain applicable to the replacement data file") } diff --git a/table/replace_files_test.go b/table/replace_files_test.go index 46f20efb0..518b5d294 100644 --- a/table/replace_files_test.go +++ b/table/replace_files_test.go @@ -20,11 +20,15 @@ package table_test import ( "context" "fmt" + "os" "path/filepath" + "slices" "strconv" + "strings" "testing" "github.com/apache/arrow-go/v18/arrow/array" + "github.com/apache/arrow-go/v18/arrow/memory" "github.com/apache/iceberg-go" iceio "github.com/apache/iceberg-go/io" "github.com/apache/iceberg-go/table" @@ -213,7 +217,7 @@ func TestReplaceFilesWithDeleteFilesPreservesDataSequenceNumber(t *testing.T) { require.NoError(t, err) oldDelete := oldDeleteBuilder.Build() - oldSequence := deleteFileSequence(t, tbl, oldDelete.FilePath()) + oldSequence := currentManifestEntry(t, tbl, oldDelete.FilePath()).SequenceNum() newDeletePath := tbl.Location() + "/data/new-pos-delete.parquet" writeParquetFile(t, newDeletePath, table.PositionalDeleteArrowSchema, fmt.Sprintf(`[{"file_path":%q,"pos":0}]`, dataPath)) @@ -232,7 +236,7 @@ func TestReplaceFilesWithDeleteFilesPreservesDataSequenceNumber(t *testing.T) { tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - assert.Equal(t, oldSequence, deleteFileSequence(t, tbl, newDeletePath), + assert.Equal(t, oldSequence, currentManifestEntry(t, tbl, newDeletePath).SequenceNum(), "rewritten delete files must retain the replaced data sequence number") assert.NotEqual(t, tbl.CurrentSnapshot().SequenceNumber, oldSequence, "the new snapshot sequence must not replace the delete data sequence") @@ -362,7 +366,7 @@ func TestReplaceFilesWithDeleteFilesValidatesDeletionVectorIdentity(t *testing.T require.NoError(t, tx.NewRowDelta(nil).AddDeletes(sourceDelete).Commit(t.Context())) sharedTbl, err := tx.Commit(t.Context()) require.NoError(t, err) - sequence := deleteFileSequence(t, sharedTbl, sourceDelete.FilePath()) + sequence := currentManifestEntry(t, sharedTbl, sourceDelete.FilePath()).SequenceNum() tx = sharedTbl.NewTransaction() require.NoError(t, tx.UpgradeFormatVersion(3)) @@ -454,7 +458,7 @@ func TestReplaceFilesWithDeleteFilesValidatesExistingPaths(t *testing.T) { } tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - sequence := deleteFileSequence(t, tbl, source.FilePath()) + sequence := currentManifestEntry(t, tbl, source.FilePath()).SequenceNum() for _, tt := range []struct { name string @@ -525,7 +529,7 @@ func TestReplaceFilesWithDeleteFilesRejectsPartialPositionDeleteToDVRewrite(t *t tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - sequence := deleteFileSequence(t, tbl, oldDeletes[0].FilePath()) + sequence := currentManifestEntry(t, tbl, oldDeletes[0].FilePath()).SequenceNum() offset, length := int64(8), int64(16) replacement := newRewriteDeletionVector(t, tbl.Location()+"/data/replacement.puffin", dataPath, &offset, &length) @@ -568,7 +572,7 @@ func TestReplaceFilesWithDeleteFilesRejectsDeletionVectorPartitionMismatch(t *te require.NoError(t, tx.NewRowDelta(nil).AddDeletes(oldDelete).Commit(t.Context())) tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - deleteSequence := fileDataSequence(t, tbl, oldDeletePath) + deleteSequence := currentManifestEntry(t, tbl, oldDeletePath).SequenceNum() tx = tbl.NewTransaction() require.NoError(t, tx.UpgradeFormatVersion(3)) @@ -620,7 +624,7 @@ func TestReplaceFilesWithDeleteFilesIgnoresOlderSurvivingPositionDelete(t *testi require.NoError(t, tx.NewRowDelta(nil).AddDeletes(oldPartitionDelete).Commit(t.Context())) tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - oldPartitionDeleteSequence := fileDataSequence(t, tbl, oldPartitionDeletePath) + oldPartitionDeleteSequence := currentManifestEntry(t, tbl, oldPartitionDeletePath).SequenceNum() dataBPath := tbl.Location() + "/data/data-b.parquet" dataBBuilder, err := iceberg.NewDataFileBuilder( @@ -631,7 +635,7 @@ func TestReplaceFilesWithDeleteFilesIgnoresOlderSurvivingPositionDelete(t *testi require.NoError(t, tx.AddDataFiles(t.Context(), []iceberg.DataFile{dataBBuilder.Build()}, nil)) tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - dataBSequence := fileDataSequence(t, tbl, dataBPath) + dataBSequence := currentManifestEntry(t, tbl, dataBPath).SequenceNum() require.Greater(t, dataBSequence, oldPartitionDeleteSequence) filePathField, ok := iceberg.PositionalDeleteSchema.FindFieldByName("file_path") @@ -651,7 +655,7 @@ func TestReplaceFilesWithDeleteFilesIgnoresOlderSurvivingPositionDelete(t *testi require.NoError(t, tx.NewRowDelta(nil).AddDeletes(newPositionDelete).Commit(t.Context())) tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - newPositionDeleteSequence := fileDataSequence(t, tbl, newPositionDeletePath) + newPositionDeleteSequence := currentManifestEntry(t, tbl, newPositionDeletePath).SequenceNum() tx = tbl.NewTransaction() require.NoError(t, tx.UpgradeFormatVersion(3)) @@ -677,7 +681,7 @@ func TestReplaceFilesWithDeleteFilesIgnoresOlderSurvivingPositionDelete(t *testi "an older position delete cannot apply to the newer target data file") tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - assert.Equal(t, newPositionDeleteSequence, fileDataSequence(t, tbl, replacement.FilePath())) + assert.Equal(t, newPositionDeleteSequence, currentManifestEntry(t, tbl, replacement.FilePath()).SequenceNum()) } func TestReplaceFilesWithDeleteFilesRejectsSurvivingDeletionVector(t *testing.T) { @@ -707,7 +711,7 @@ func TestReplaceFilesWithDeleteFilesRejectsSurvivingDeletionVector(t *testing.T) offset, length := int64(8), int64(16) replacement := newRewriteDeletionVector(t, tbl.Location()+"/data/replacement-dv.puffin", target, &offset, &length) - sequence := deleteFileSequence(t, tbl, siblingDVs[0].FilePath()) + sequence := currentManifestEntry(t, tbl, siblingDVs[0].FilePath()).SequenceNum() tx = tbl.NewTransaction() err = tx.ReplaceFilesWithDeleteFiles(t.Context(), nil, nil, []iceberg.DataFile{siblingDVs[0]}, @@ -729,7 +733,7 @@ func TestReplaceFilesWithDeleteFilesAllowsDroppedEqualityField(t *testing.T) { require.NoError(t, tx.NewRowDelta(nil).AddDeletes(oldDelete).Commit(t.Context())) tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - sequence := deleteFileSequence(t, tbl, oldDelete.FilePath()) + sequence := currentManifestEntry(t, tbl, oldDelete.FilePath()).SequenceNum() tx = tbl.NewTransaction() require.NoError(t, tx.UpdateSchema(true, false).DeleteColumn([]string{"data"}).Commit()) @@ -750,11 +754,9 @@ func TestReplaceFilesWithDeleteFilesAllowsDroppedEqualityField(t *testing.T) { []table.DeleteFileAddition{{File: newDelete, DataSequenceNumber: sequence}}, nil)) } -func deleteFileSequence(t *testing.T, tbl *table.Table, path string) int64 { - return fileDataSequence(t, tbl, path) -} - -func fileDataSequence(t *testing.T, tbl *table.Table, path string) int64 { +// currentManifestEntry returns the manifest entry for path in the current +// snapshot, including DELETED entries. +func currentManifestEntry(t *testing.T, tbl *table.Table, path string) iceberg.ManifestEntry { t.Helper() snap := tbl.CurrentSnapshot() require.NotNil(t, snap) @@ -764,13 +766,13 @@ func fileDataSequence(t *testing.T, tbl *table.Table, path string) int64 { for entry, err := range manifest.Entries(iceio.LocalFS{}, false) { require.NoError(t, err) if entry.DataFile().FilePath() == path { - return entry.SequenceNum() + return entry } } } t.Fatalf("file %q not found in current snapshot", path) - return -1 + return nil } func scanIDs(t *testing.T, tbl *table.Table) []int64 { @@ -936,3 +938,399 @@ func TestReplaceFiles_ValidationErrors(t *testing.T) { assert.Contains(t, err.Error(), "cannot remove deletion vectors that do not belong to the table") }) } + +// seedDataFilesWithPositionDeletes creates two v2 data files (ids 1-3 and +// 4-6), each with a position delete on a different row. Tasks are sorted by path. +func seedDataFilesWithPositionDeletes(t *testing.T) (*table.Table, []table.FileScanTask) { + t.Helper() + + tbl := newReplaceFilesTestTable(t) + arrowSc, err := table.SchemaToArrowSchema(tbl.Schema(), nil, false, false) + require.NoError(t, err) + + files := []struct { + rows string + deletePos int64 + }{ + { + rows: `[{"id":1,"data":"a"}, {"id":2,"data":"b"}, {"id":3,"data":"c"}]`, + deletePos: 0, + }, + { + rows: `[{"id":4,"data":"d"}, {"id":5,"data":"e"}, {"id":6,"data":"f"}]`, + deletePos: 1, + }, + } + dataPaths := make([]string, len(files)) + for i, f := range files { + dataPaths[i] = fmt.Sprintf("%s/data/data-%d.parquet", tbl.Location(), i) + writeParquetFile(t, dataPaths[i], arrowSc, f.rows) + } + + tx := tbl.NewTransaction() + require.NoError(t, tx.AddFiles(t.Context(), dataPaths, nil, false)) + tbl, err = tx.Commit(t.Context()) + require.NoError(t, err) + + filePathField, ok := iceberg.PositionalDeleteSchema.FindFieldByName("file_path") + require.True(t, ok) + tx = tbl.NewTransaction() + rowDelta := tx.NewRowDelta(nil) + for i, f := range files { + deletePath := fmt.Sprintf("%s/data/pos-delete-%d.parquet", tbl.Location(), i) + row := fmt.Sprintf(`[{"file_path":%q,"pos":%d}]`, dataPaths[i], f.deletePos) + + writeParquetFile(t, deletePath, table.PositionalDeleteArrowSchema, row) + bound, err := iceberg.StringLiteral(dataPaths[i]).MarshalBinary() + require.NoError(t, err) + + builder, err := iceberg.NewDataFileBuilder( + *iceberg.UnpartitionedSpec, iceberg.EntryContentPosDeletes, + deletePath, iceberg.ParquetFile, nil, nil, nil, 1, 128, + ) + require.NoError(t, err) + rowDelta.AddDeletes(builder. + LowerBoundValues(map[int][]byte{filePathField.ID: bound}). + UpperBoundValues(map[int][]byte{filePathField.ID: bound}). + Build()) + } + + require.NoError(t, rowDelta.Commit(t.Context())) + tbl, err = tx.Commit(t.Context()) + require.NoError(t, err) + require.Equal(t, []int64{2, 3, 4, 6}, idsInTable(t, tbl)) + + tasks, err := tbl.Scan().PlanFiles(t.Context()) + require.NoError(t, err) + require.Len(t, tasks, len(files)) + slices.SortFunc(tasks, func(a, b table.FileScanTask) int { + return strings.Compare(a.File.FilePath(), b.File.FilePath()) + }) + for _, task := range tasks { + require.Len(t, task.DeleteFiles, 1, "each data file must carry only its own position delete") + } + + return tbl, tasks +} + +// seedReplaceFilesTableWithDelete returns a table with one data file and one +// delete on it: a position delete on v2, a deletion vector on v3. +func seedReplaceFilesTableWithDelete(t *testing.T, version int) (*table.Table, iceberg.DataFile, iceberg.DataFile) { + t.Helper() + + switch version { + case 2: + tbl, tasks := seedDataFilesWithPositionDeletes(t) + + return tbl, tasks[0].File, tasks[0].DeleteFiles[0] + case 3: + tbl, target := seedV3TableWithDV(t) + tasks, err := tbl.Scan().PlanFiles(t.Context()) + require.NoError(t, err) + + for _, task := range tasks { + if task.File.FilePath() == target { + require.Len(t, task.DeletionVectorFiles, 1) + + return tbl, task.File, task.DeletionVectorFiles[0] + } + } + t.Fatalf("deletion vector target %q not found in scan tasks", target) + default: + t.Fatalf("unsupported format version %d", version) + } + + return nil, nil, nil +} + +// seedSupersededDeletionVector deletes twice from one v3 data file. The first +// DV gets replaced and stays only as a DELETED entry. It returns that old DV +// and the scan task holding the new, live DV. +func seedSupersededDeletionVector(t *testing.T) (*table.Table, table.FileScanTask, iceberg.DataFile) { + t.Helper() + + tbl := newMergeOnReadTestTableVersion(t, "3") + arrowSc, err := table.SchemaToArrowSchema(tbl.Schema(), nil, false, false) + require.NoError(t, err) + + data, err := array.TableFromJSON(memory.DefaultAllocator, arrowSc, []string{ + `[{"id":1,"data":"a"},{"id":2,"data":"b"}, + {"id":3,"data":"c"},{"id":4,"data":"d"}, + {"id":5,"data":"e"}]`, + }) + require.NoError(t, err) + defer data.Release() + + tbl, err = tbl.Append(t.Context(), array.NewTableReader(data, -1), nil) + require.NoError(t, err) + + var superseded iceberg.DataFile + for _, id := range []int64{2, 4} { + tbl, err = tbl.Delete(t.Context(), iceberg.EqualTo(iceberg.Reference("id"), id), nil) + require.NoError(t, err) + + if superseded == nil { + tasks, err := tbl.Scan().PlanFiles(t.Context()) + require.NoError(t, err) + require.Len(t, tasks, 1) + require.Len(t, tasks[0].DeletionVectorFiles, 1) + superseded = tasks[0].DeletionVectorFiles[0] + } + } + require.Equal(t, 1, liveDVCount(t, tbl)) + require.Equal(t, iceberg.EntryStatusDELETED, currentManifestEntry(t, tbl, superseded.FilePath()).Status(), + "the second delete must keep the first DV only as a DELETED entry") + + tasks, err := tbl.Scan().PlanFiles(t.Context()) + require.NoError(t, err) + require.Len(t, tasks, 1) + require.Len(t, tasks[0].DeletionVectorFiles, 1) + require.Equal(t, superseded.ReferencedDataFile(), tasks[0].DeletionVectorFiles[0].ReferencedDataFile()) + require.NotEqual(t, superseded.FilePath(), tasks[0].DeletionVectorFiles[0].FilePath()) + + return tbl, tasks[0], superseded +} + +func newCompactedDataFile(t *testing.T, tbl *table.Table, path string, recordCount int64) iceberg.DataFile { + t.Helper() + + builder, err := iceberg.NewDataFileBuilder( + *iceberg.UnpartitionedSpec, iceberg.EntryContentData, + path, iceberg.ParquetFile, nil, nil, nil, recordCount, 512, + ) + require.NoError(t, err) + if tbl.Metadata().Version() >= 3 { + builder.FirstRowID(tbl.Metadata().NextRowID()) + } + + return builder.Build() +} + +// compactWithDeletes swaps data for compacted and removes dels, either through +// ReplaceFiles or, when automatic is set, through RewriteFiles. +func compactWithDeletes(ctx context.Context, tx *table.Transaction, automatic bool, data, compacted, dels []iceberg.DataFile) error { + if !automatic { + return tx.ReplaceFiles(ctx, data, compacted, dels, nil) + } + + result := table.CompactionGroupResult{OldDataFiles: data, NewDataFiles: compacted} + for _, del := range dels { + if table.IsDeletionVector(del) { + result.SafeDeletionVectors = append(result.SafeDeletionVectors, del) + } else { + result.SafePosDeletes = append(result.SafePosDeletes, del) + } + } + + return tx.NewRewrite(nil).ApplyResult(result).Commit(ctx) +} + +// metadataFileCount counts files in the table's metadata folder. Staging a +// commit writes files there, so an unchanged count means nothing was written. +func metadataFileCount(t *testing.T, tbl *table.Table) int { + t.Helper() + + entries, err := os.ReadDir(filepath.FromSlash(tbl.Location() + "/metadata")) + require.NoError(t, err) + + return len(entries) +} + +// A compaction planned before another writer deleted one of its data files +// must fail. Committing it would bring the deleted rows back. +func TestReplaceFilesRejectsDataFileDeletedByCurrentSnapshot(t *testing.T) { + tbl, tasks := seedDataFilesWithPositionDeletes(t) + arrowSc, err := table.SchemaToArrowSchema(tbl.Schema(), nil, false, false) + require.NoError(t, err) + compactedPath := tbl.Location() + "/data/compacted-0.parquet" + + jsonData := `[{"id":2,"data":"b"},{"id":3,"data":"c"}]` + writeParquetFile(t, compactedPath, arrowSc, jsonData) + + path := tbl.Location() + "/data/compacted-1.parquet" + compacted := []iceberg.DataFile{ + newCompactedDataFile(t, tbl, compactedPath, 2), + newCompactedDataFile(t, tbl, path, 2), + } + + tx := tbl.NewTransaction() + require.NoError(t, tx.Delete(t.Context(), iceberg.GreaterThanEqual(iceberg.Reference("id"), int64(4)), nil)) + tbl, err = tx.Commit(t.Context()) + require.NoError(t, err) + require.Equal(t, []int64{2, 3}, idsInTable(t, tbl)) + require.Equal(t, iceberg.EntryStatusDELETED, currentManifestEntry(t, tbl, tasks[1].File.FilePath()).Status()) + require.NotEqual(t, iceberg.EntryStatusDELETED, currentManifestEntry(t, tbl, tasks[0].File.FilePath()).Status()) + + before := metadataFileCount(t, tbl) + tx = tbl.NewTransaction() + err = tx.ReplaceFiles(t.Context(), + []iceberg.DataFile{tasks[0].File, tasks[1].File}, compacted, + []iceberg.DataFile{tasks[0].DeleteFiles[0], tasks[1].DeleteFiles[0]}, nil) + assert.ErrorContains(t, err, "cannot delete data files that do not belong to the table") + assert.Equal(t, before, metadataFileCount(t, tbl), "a rejected replace must not write manifests") + + tx = tbl.NewTransaction() + require.NoError(t, tx.ReplaceFiles(t.Context(), []iceberg.DataFile{tasks[0].File}, compacted[:1], tasks[0].DeleteFiles, nil), + "the data file the delete left live must still compact") + assert.Greater(t, metadataFileCount(t, tbl), before, "staging a replace must write manifests") + + tbl, err = tx.Commit(t.Context()) + require.NoError(t, err) + assert.Equal(t, []int64{2, 3}, idsInTable(t, tbl)) +} + +func TestReplaceFilesRejectsDeleteFileRemovedByCurrentSnapshot(t *testing.T) { + for _, tt := range []struct { + name string + version int + automatic bool + wantErr string + }{ + { + name: "position delete", + version: 2, + automatic: false, + wantErr: "cannot remove delete files that do not belong to the table", + }, + { + name: "automatic position delete", + version: 2, + automatic: true, + wantErr: "cannot remove automatic delete files that do not belong to the table", + }, + { + name: "deletion vector", + version: 3, + automatic: false, + wantErr: "cannot remove deletion vectors that do not belong to the table", + }, + { + name: "automatic deletion vector", + version: 3, + automatic: true, + wantErr: "cannot remove deletion vectors that do not belong to the table", + }, + } { + t.Run(tt.name, func(t *testing.T) { + tbl, data, del := seedReplaceFilesTableWithDelete(t, tt.version) + + tx := tbl.NewTransaction() + require.NoError(t, tx.ReplaceFiles(t.Context(), nil, nil, []iceberg.DataFile{del}, nil)) + tbl, err := tx.Commit(t.Context()) + require.NoError(t, err) + require.Equal(t, iceberg.EntryStatusDELETED, currentManifestEntry(t, tbl, del.FilePath()).Status()) + require.NotEqual(t, iceberg.EntryStatusDELETED, currentManifestEntry(t, tbl, data.FilePath()).Status()) + + compacted := newCompactedDataFile(t, tbl, tbl.Location()+"/data/compacted.parquet", 1) + before := metadataFileCount(t, tbl) + err = compactWithDeletes(t.Context(), tbl.NewTransaction(), tt.automatic, + []iceberg.DataFile{data}, []iceberg.DataFile{compacted}, []iceberg.DataFile{del}) + assert.ErrorContains(t, err, tt.wantErr) + assert.Equal(t, before, metadataFileCount(t, tbl), "a rejected replace must not write manifests") + }) + } +} + +// A rewrite planned with an old DV must fail once a newer DV has replaced it, +// because the newer DV has deletes the rewrite never saw. +func TestReplaceFilesRejectsDeletionVectorSupersededByCurrentSnapshot(t *testing.T) { + type stageFunc func(t *testing.T, tx *table.Transaction, tbl *table.Table, task table.FileScanTask, superseded iceberg.DataFile) error + compaction := func(automatic bool) stageFunc { + return func(t *testing.T, tx *table.Transaction, tbl *table.Table, task table.FileScanTask, superseded iceberg.DataFile) error { + compacted := newCompactedDataFile(t, tbl, tbl.Location()+"/data/compacted.parquet", 4) + + return compactWithDeletes(t.Context(), tx, automatic, + []iceberg.DataFile{task.File}, []iceberg.DataFile{compacted}, []iceberg.DataFile{superseded}) + } + } + + for _, tt := range []struct { + name string + stage stageFunc + }{ + { + name: "compaction", + stage: compaction(false), + }, + { + name: "automatic compaction", + stage: compaction(true), + }, + { + name: "deletion vector rewrite", + stage: func(t *testing.T, tx *table.Transaction, tbl *table.Table, task table.FileScanTask, superseded iceberg.DataFile) error { + writer := dv.NewDVWriter(iceio.LocalFS{}, unpartitionedSpecByID) + require.NoError(t, writer.Add(task.File.FilePath(), []int64{1}, 0, nil)) + rewritten, err := writer.Flush(t.Context(), tbl.Location()+"/data/rewritten-dv.puffin") + require.NoError(t, err) + require.Len(t, rewritten, 1) + + return tx.ReplaceFilesWithDeleteFiles(t.Context(), nil, nil, + []iceberg.DataFile{superseded}, + []table.DeleteFileAddition{{ + File: rewritten[0], + DataSequenceNumber: currentManifestEntry(t, tbl, superseded.FilePath()).SequenceNum(), + }}, nil) + }, + }, + } { + t.Run(tt.name, func(t *testing.T) { + tbl, task, superseded := seedSupersededDeletionVector(t) + + before := metadataFileCount(t, tbl) + err := tt.stage(t, tbl.NewTransaction(), tbl, task, superseded) + assert.ErrorContains(t, err, "cannot remove deletion vectors that do not belong to the table") + assert.Equal(t, before, metadataFileCount(t, tbl), "a rejected replace must not write manifests") + }) + } +} + +// Removing the current DV must remove that DV, not the old replaced one that +// is still listed as a DELETED entry. +func TestReplaceFilesRemovesLiveDeletionVectorNotSupersededEntry(t *testing.T) { + t.Run("compaction", func(t *testing.T) { + tbl, task, _ := seedSupersededDeletionVector(t) + arrowSc, err := table.SchemaToArrowSchema(tbl.Schema(), nil, false, false) + require.NoError(t, err) + + compactedPath := tbl.Location() + "/data/compacted.parquet" + jsonData := `[{"id":1,"data":"a"},{"id":3,"data":"c"},{"id":5,"data":"e"}]` + writeParquetFile(t, compactedPath, arrowSc, jsonData) + + tx := tbl.NewTransaction() + require.NoError(t, tx.ReplaceFiles(t.Context(), + []iceberg.DataFile{task.File}, + []iceberg.DataFile{newCompactedDataFile(t, tbl, compactedPath, 3)}, + task.DeletionVectorFiles, nil)) + tbl, err = tx.Commit(t.Context()) + require.NoError(t, err) + + assert.Zero(t, liveDVCount(t, tbl), "the live deletion vector must be removed with its data file") + assert.Equal(t, []int64{1, 3, 5}, idsInTable(t, tbl)) + }) + + t.Run("deletion vector rewrite", func(t *testing.T) { + tbl, task, _ := seedSupersededDeletionVector(t) + liveDV := task.DeletionVectorFiles[0] + writer := dv.NewDVWriter(iceio.LocalFS{}, unpartitionedSpecByID) + require.NoError(t, writer.Add(task.File.FilePath(), []int64{1, 3}, 0, nil)) + + rewritten, err := writer.Flush(t.Context(), tbl.Location()+"/data/rewritten-dv.puffin") + require.NoError(t, err) + require.Len(t, rewritten, 1) + + tx := tbl.NewTransaction() + require.NoError(t, tx.ReplaceFilesWithDeleteFiles(t.Context(), nil, nil, + []iceberg.DataFile{liveDV}, + []table.DeleteFileAddition{{ + File: rewritten[0], + DataSequenceNumber: currentManifestEntry(t, tbl, liveDV.FilePath()).SequenceNum(), + }}, nil)) + tbl, err = tx.Commit(t.Context()) + require.NoError(t, err) + + assert.Equal(t, 1, liveDVCount(t, tbl), "a data file must keep exactly one deletion vector") + assert.Equal(t, iceberg.EntryStatusDELETED, currentManifestEntry(t, tbl, liveDV.FilePath()).Status()) + assert.Equal(t, []int64{1, 3, 5}, idsInTable(t, tbl)) + }) +} diff --git a/table/rewrite_files_test.go b/table/rewrite_files_test.go index d205d62ae..63ed00b25 100644 --- a/table/rewrite_files_test.go +++ b/table/rewrite_files_test.go @@ -327,7 +327,7 @@ func TestRewriteFiles_PreservesIndependentPositionDeleteSequences(t *testing.T) tbl, err = tx.Commit(t.Context()) require.NoError(t, err) oldDeletes = append(oldDeletes, builder) - oldSequences = append(oldSequences, deleteFileSequence(t, tbl, path)) + oldSequences = append(oldSequences, currentManifestEntry(t, tbl, path).SequenceNum()) } require.NotEqual(t, oldSequences[0], oldSequences[1]) @@ -349,7 +349,7 @@ func TestRewriteFiles_PreservesIndependentPositionDeleteSequences(t *testing.T) require.NoError(t, err) for i, df := range newDeletes { - assert.Equal(t, oldSequences[i], deleteFileSequence(t, tbl, df.FilePath())) + assert.Equal(t, oldSequences[i], currentManifestEntry(t, tbl, df.FilePath()).SequenceNum()) } assert.Empty(t, scanIDs(t, tbl), "each replacement must retain the applicability of its corresponding old position delete") @@ -814,7 +814,7 @@ func TestRewriteFiles_RejectsConcurrentDeleteForAddedDVTarget(t *testing.T) { require.NoError(t, tx.NewRowDelta(nil).AddDeletes(oldDelete).Commit(t.Context())) tbl, err = tx.Commit(t.Context()) require.NoError(t, err) - oldDeleteSequence := fileDataSequence(t, tbl, oldDelete.FilePath()) + oldDeleteSequence := currentManifestEntry(t, tbl, oldDelete.FilePath()).SequenceNum() tx = tbl.NewTransaction() require.NoError(t, tx.UpgradeFormatVersion(3)) diff --git a/table/row_delta.go b/table/row_delta.go index 0e9508528..a8fee000f 100644 --- a/table/row_delta.go +++ b/table/row_delta.go @@ -435,13 +435,10 @@ func (rd *RowDelta) resolveRemovedDeletes(fs iceio.IO, meta *MetadataBuilder) (r } liveByPath := make(map[string][]iceberg.DataFile, len(rd.removedDels)) - for entry, err := range snap.entries(fs, iceberg.ManifestContentDeletes) { + for entry, err := range snap.entries(fs, iceberg.ManifestContentDeletes, true) { if err != nil { return nil, nil, err } - if entry.Status() == iceberg.EntryStatusDELETED { - continue - } df := entry.DataFile() if _, ok := want[df.FilePath()]; ok { liveByPath[df.FilePath()] = append(liveByPath[df.FilePath()], df) diff --git a/table/snapshots.go b/table/snapshots.go index 8869515c1..4d084332b 100644 --- a/table/snapshots.go +++ b/table/snapshots.go @@ -508,8 +508,9 @@ func (s Snapshot) dataFiles(fio iceio.IO, fileFilter set[iceberg.ManifestEntryCo // matters. // // manifestContent < 0 yields entries across both data and delete -// manifests. -func (s Snapshot) entries(fio iceio.IO, manifestContent iceberg.ManifestContent) iter.Seq2[iceberg.ManifestEntry, error] { +// manifests. Set discardDeleted to skip DELETED entries: they are files +// already removed from the table, not files still in it. +func (s Snapshot) entries(fio iceio.IO, manifestContent iceberg.ManifestContent, discardDeleted bool) iter.Seq2[iceberg.ManifestEntry, error] { return func(yield func(iceberg.ManifestEntry, error) bool) { manifests, err := s.Manifests(fio) if err != nil { @@ -523,7 +524,7 @@ func (s Snapshot) entries(fio iceio.IO, manifestContent iceberg.ManifestContent) continue } - for entry, err := range m.Entries(fio, false) { + for entry, err := range m.Entries(fio, discardDeleted) { if err != nil { yield(nil, err) diff --git a/table/transaction.go b/table/transaction.go index ab03f560a..22687295f 100644 --- a/table/transaction.go +++ b/table/transaction.go @@ -1846,7 +1846,7 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF } setDeleteFilesToRemove := make(map[string]struct{}, len(deleteFilesToRemove)) - dvRefsToRemove := make(map[string]struct{}, len(deleteFilesToRemove)) + dvsToRemoveByRef := make(map[string]iceberg.DataFile, len(deleteFilesToRemove)) for i, df := range deleteFilesToRemove { if df == nil { return fmt.Errorf("nil delete file at index %d for ReplaceFiles", i) @@ -1860,10 +1860,10 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF if ref == nil { return errors.New("deletion vector to remove is missing referenced_data_file for ReplaceFiles") } - if _, ok := dvRefsToRemove[*ref]; ok { + if _, ok := dvsToRemoveByRef[*ref]; ok { return errors.New("deletion vectors to remove must reference distinct data files for ReplaceFiles") } - dvRefsToRemove[*ref] = struct{}{} + dvsToRemoveByRef[*ref] = df continue } @@ -1873,7 +1873,7 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF setDeleteFilesToRemove[path] = struct{}{} } autoSetDeleteFilesToRemove := make(map[string]struct{}, len(autoDeleteFilesToRemove)) - autoDVRefsToRemove := make(map[string]struct{}, len(autoDeleteFilesToRemove)) + autoDVsToRemoveByRef := make(map[string]iceberg.DataFile, len(autoDeleteFilesToRemove)) for i, df := range autoDeleteFilesToRemove { if df == nil { return fmt.Errorf("nil automatic delete file at index %d for ReplaceFiles", i) @@ -1887,13 +1887,13 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF if ref == nil { return errors.New("automatic deletion vector to remove is missing referenced_data_file for ReplaceFiles") } - if _, ok := dvRefsToRemove[*ref]; ok { + if _, ok := dvsToRemoveByRef[*ref]; ok { return errors.New("delete vectors to remove must reference distinct data files for ReplaceFiles") } - if _, ok := autoDVRefsToRemove[*ref]; ok { + if _, ok := autoDVsToRemoveByRef[*ref]; ok { return errors.New("automatic deletion vectors to remove must reference distinct data files for ReplaceFiles") } - autoDVRefsToRemove[*ref] = struct{}{} + autoDVsToRemoveByRef[*ref] = df continue } @@ -1925,26 +1925,28 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF markedDataForDeletion := make([]iceberg.DataFile, 0, len(setToDelete)) markedDeleteForRemoval := make([]iceberg.DataFile, 0, len(setDeleteFilesToRemove)) markedAutoDeleteForRemoval := make([]iceberg.DataFile, 0, len(autoSetDeleteFilesToRemove)) - markedDVsForRemoval := make(map[string]iceberg.DataFile, len(dvRefsToRemove)) + markedDVsForRemoval := make(map[string]iceberg.DataFile, len(dvsToRemoveByRef)) removedDeleteSequenceNumbers := make([]int64, 0, len(deleteFilesToRemove)) removedDeleteContents := make(map[iceberg.ManifestEntryContent]struct{}) liveDataFiles := make(map[string]rewriteFileState) survivingPositionDeletes := make([]rewriteFileState, 0) - for entry, err := range s.entries(fs, -1) { + for entry, err := range s.entries(fs, -1, false) { if err != nil { return err } df := entry.DataFile() path := df.FilePath() isData := df.ContentType() == iceberg.EntryContentData + // A DELETED entry is a file already removed from the table. It must not + // count as found when removing files, but it still blocks re-adding its path. isLive := entry.Status() != iceberg.EntryStatusDELETED if isData && isLive { liveDataFiles[path] = rewriteFileState{file: df, dataSequenceNumber: entry.SequenceNum()} + if _, ok := setToDelete[path]; ok { + markedDataForDeletion = append(markedDataForDeletion, df) + } } - if _, ok := setToDelete[path]; ok && isData { - markedDataForDeletion = append(markedDataForDeletion, df) - } - if !isData { + if !isData && isLive { if _, ok := setDeleteFilesToRemove[path]; ok { markedDeleteForRemoval = append(markedDeleteForRemoval, df) if seq := entry.SequenceNum(); seq >= 0 { @@ -1956,7 +1958,9 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF } else if _, ok := autoSetDeleteFilesToRemove[path]; ok { markedAutoDeleteForRemoval = append(markedAutoDeleteForRemoval, df) } else if ref := iceberginternal.BorrowedDataFileReferencedDataFile(df); IsDeletionVector(df) && ref != nil { - if _, ok := dvRefsToRemove[*ref]; ok { + // Match the DV's path too, not only its data file. Otherwise an old DV that + // was already replaced would match the newer DV, which holds more deletes. + if want, ok := dvsToRemoveByRef[*ref]; ok && want.FilePath() == path { markedDVsForRemoval[*ref] = df if seq := entry.SequenceNum(); seq >= 0 { removedDeleteSequenceNumbers = append(removedDeleteSequenceNumbers, seq) @@ -1964,7 +1968,7 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF } else { return fmt.Errorf("deletion vector %s has no data sequence number in the current snapshot", path) } - } else if _, ok := autoDVRefsToRemove[*ref]; ok { + } else if want, ok := autoDVsToRemoveByRef[*ref]; ok && want.FilePath() == path { markedDVsForRemoval[*ref] = df } } @@ -2001,10 +2005,10 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF if _, addingReplacement := setDeleteFilesToAdd.dvsByRef[*ref]; !addingReplacement { continue } - if _, explicitlyRemoved := dvRefsToRemove[*ref]; explicitlyRemoved { + if _, explicitlyRemoved := dvsToRemoveByRef[*ref]; explicitlyRemoved { continue } - if _, automaticallyRemoved := autoDVRefsToRemove[*ref]; automaticallyRemoved { + if _, automaticallyRemoved := autoDVsToRemoveByRef[*ref]; automaticallyRemoved { continue } @@ -2043,9 +2047,9 @@ func (t *Transaction) replaceFiles(ctx context.Context, dataFilesToDelete, dataF if len(markedAutoDeleteForRemoval) != len(autoSetDeleteFilesToRemove) { return errors.New("cannot remove automatic delete files that do not belong to the table") } - // Keyed by referenced data file, so duplicate DV entries for one ref collapse - // to one slot; equality then means every requested ref exists in the table. - if len(markedDVsForRemoval) != len(dvRefsToRemove)+len(autoDVRefsToRemove) { + // markedDVsForRemoval has one slot per data file, so equal counts mean + // every requested DV was found live in the table. + if len(markedDVsForRemoval) != len(dvsToRemoveByRef)+len(autoDVsToRemoveByRef) { return errors.New("cannot remove deletion vectors that do not belong to the table") } @@ -3142,19 +3146,12 @@ func (t *Transaction) collectExistingDVs(fs io.IO, files []iceberg.DataFile) (ma } result := make(map[string]iceberg.DataFile) - // Iterate delete manifests only and skip DELETED-status entries: a - // superseded DV lingers as a DELETED entry in the deleted-files manifest - // against the same referenced data file. Including it here would let the - // stale ghost win the last-write into result (manifest concat order places - // deleted entries last), seeding the new DV from an outdated bitmap and - // resurrecting rows removed by the prior delete. See issue #1372. - for entry, err := range s.entries(fs, iceberg.ManifestContentDeletes) { + // Skip DELETED entries: an old, replaced DV stays behind as one. Reading it + // would build the new DV from old data and bring back deleted rows. See #1372. + for entry, err := range s.entries(fs, iceberg.ManifestContentDeletes, true) { if err != nil { return nil, fmt.Errorf("scanning existing deletion vectors: %w", err) } - if entry.Status() == iceberg.EntryStatusDELETED { - continue - } df := entry.DataFile() if !IsDeletionVector(df) { continue diff --git a/table/transaction_internal_test.go b/table/transaction_internal_test.go index d673b2d43..7754928c3 100644 --- a/table/transaction_internal_test.go +++ b/table/transaction_internal_test.go @@ -174,11 +174,8 @@ func liveDataFilePathsForSnapshot(t *testing.T, snap *Snapshot, fs iceio.IO) []s t.Helper() require.NotNil(t, snap) var paths []string - for e, err := range snap.entries(fs, iceberg.ManifestContentData) { + for e, err := range snap.entries(fs, iceberg.ManifestContentData, true) { require.NoError(t, err) - if e.Status() == iceberg.EntryStatusDELETED { - continue - } if e.DataFile().ContentType() == iceberg.EntryContentData { paths = append(paths, e.DataFile().FilePath()) } From 0e6884c75d1d024d6cb0e06cafc3784a878b01e2 Mon Sep 17 00:00:00 2001 From: badalprasadsingh Date: Fri, 25 Sep 2026 18:41:26 +0530 Subject: [PATCH 2/2] chore: minor Signed-off-by: badalprasadsingh --- table/replace_files_test.go | 22 +++++++++++----------- 1 file changed, 11 insertions(+), 11 deletions(-) diff --git a/table/replace_files_test.go b/table/replace_files_test.go index 518b5d294..a8e2e62fc 100644 --- a/table/replace_files_test.go +++ b/table/replace_files_test.go @@ -939,8 +939,8 @@ func TestReplaceFiles_ValidationErrors(t *testing.T) { }) } -// seedDataFilesWithPositionDeletes creates two v2 data files (ids 1-3 and -// 4-6), each with a position delete on a different row. Tasks are sorted by path. +// seedDataFilesWithPositionDeletes creates two v2 data files (ids 1-3 and 4-6) +// each with a position delete on a different row. Tasks are sorted by path. func seedDataFilesWithPositionDeletes(t *testing.T) (*table.Table, []table.FileScanTask) { t.Helper() @@ -1013,8 +1013,8 @@ func seedDataFilesWithPositionDeletes(t *testing.T) (*table.Table, []table.FileS return tbl, tasks } -// seedReplaceFilesTableWithDelete returns a table with one data file and one -// delete on it: a position delete on v2, a deletion vector on v3. +// seedReplaceFilesTableWithDelete returns a table with one data file and one delete on it: +// a position delete on v2, a deletion vector on v3. func seedReplaceFilesTableWithDelete(t *testing.T, version int) (*table.Table, iceberg.DataFile, iceberg.DataFile) { t.Helper() @@ -1043,9 +1043,9 @@ func seedReplaceFilesTableWithDelete(t *testing.T, version int) (*table.Table, i return nil, nil, nil } -// seedSupersededDeletionVector deletes twice from one v3 data file. The first -// DV gets replaced and stays only as a DELETED entry. It returns that old DV -// and the scan task holding the new, live DV. +// seedSupersededDeletionVector deletes twice from one v3 data file. +// The first DV gets replaced and stays only as a DELETED entry. +// It returns that old DV and the scan task holding the new, live DV. func seedSupersededDeletionVector(t *testing.T) (*table.Table, table.FileScanTask, iceberg.DataFile) { t.Helper() @@ -1125,8 +1125,8 @@ func compactWithDeletes(ctx context.Context, tx *table.Transaction, automatic bo return tx.NewRewrite(nil).ApplyResult(result).Commit(ctx) } -// metadataFileCount counts files in the table's metadata folder. Staging a -// commit writes files there, so an unchanged count means nothing was written. +// metadataFileCount counts files in the table's metadata folder. +// Staging a commit writes files there, so an unchanged count means nothing was written. func metadataFileCount(t *testing.T, tbl *table.Table) int { t.Helper() @@ -1136,8 +1136,8 @@ func metadataFileCount(t *testing.T, tbl *table.Table) int { return len(entries) } -// A compaction planned before another writer deleted one of its data files -// must fail. Committing it would bring the deleted rows back. +// A compaction planned before another writer deleted one of its data files must fail. +// Committing it would bring the deleted rows back. func TestReplaceFilesRejectsDataFileDeletedByCurrentSnapshot(t *testing.T) { tbl, tasks := seedDataFilesWithPositionDeletes(t) arrowSc, err := table.SchemaToArrowSchema(tbl.Schema(), nil, false, false)