From 88ebdc042bf4a78f1390d6bf81b520012f63dd4c Mon Sep 17 00:00:00 2001 From: Digvijay Date: Wed, 26 Aug 2026 23:57:15 -0500 Subject: [PATCH 1/3] feat(table): add delete-file-backed changelog task types Introduce AddedRowsScanTask, DeletedDataFileScanTask, and DeletedRowsScanTask so later changelog planning can distinguish row inserts from file removal and added versus existing deletes. Signed-off-by: Digvijay --- table/changelog_scan_task.go | 123 ++++++++++++++++++++++++++++++ table/changelog_scan_task_test.go | 72 +++++++++++++++++ 2 files changed, 195 insertions(+) create mode 100644 table/changelog_scan_task.go create mode 100644 table/changelog_scan_task_test.go diff --git a/table/changelog_scan_task.go b/table/changelog_scan_task.go new file mode 100644 index 000000000..279c368e6 --- /dev/null +++ b/table/changelog_scan_task.go @@ -0,0 +1,123 @@ +// 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 table + +import "github.com/apache/iceberg-go" + +// AddedRowsScanTask is a changelog insert produced by adding a data file. +// Matching delete files committed in the same snapshot, or from squashed +// snapshots, are applied while reading so deleted rows are not emitted as +// inserts. +type AddedRowsScanTask struct { + FileScanTask + changeOrdinal int + commitSnapshotID int64 +} + +// NewAddedRowsScanTask constructs an insert task for dataFile. deletes are +// delete files that apply while reading the added file. +func NewAddedRowsScanTask(dataFile iceberg.DataFile, deletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) AddedRowsScanTask { + return AddedRowsScanTask{ + FileScanTask: FileScanTask{ + File: dataFile, + DeleteFiles: deletes, + }, + changeOrdinal: changeOrdinal, + commitSnapshotID: commitSnapshotID, + } +} + +func (t AddedRowsScanTask) ChangeOrdinal() int { return t.changeOrdinal } +func (t AddedRowsScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } + +// Deletes returns delete files to apply when reading the added data file. +func (t AddedRowsScanTask) Deletes() []iceberg.DataFile { + return t.DeleteFiles +} + +// DeletedDataFileScanTask is a changelog delete produced by removing a data +// file. ExistingDeletes are delete files that were already present and must +// be applied so only rows that were live when the file was removed appear as +// deletes. +type DeletedDataFileScanTask struct { + FileScanTask + changeOrdinal int + commitSnapshotID int64 +} + +// NewDeletedDataFileScanTask constructs a delete task for a removed data file. +func NewDeletedDataFileScanTask(dataFile iceberg.DataFile, existingDeletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) DeletedDataFileScanTask { + return DeletedDataFileScanTask{ + FileScanTask: FileScanTask{ + File: dataFile, + DeleteFiles: existingDeletes, + }, + changeOrdinal: changeOrdinal, + commitSnapshotID: commitSnapshotID, + } +} + +func (t DeletedDataFileScanTask) ChangeOrdinal() int { return t.changeOrdinal } +func (t DeletedDataFileScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } + +// ExistingDeletes returns delete files that applied before the data file was +// removed. +func (t DeletedDataFileScanTask) ExistingDeletes() []iceberg.DataFile { + return t.DeleteFiles +} + +// DeletedRowsScanTask is a changelog delete produced by adding delete files +// against a data file that remains in the table. AddedDeletes remove rows +// that should appear in the changelog. ExistingDeletes already applied and +// those rows must not be emitted again. +type DeletedRowsScanTask struct { + FileScanTask + addedDeletes []iceberg.DataFile + changeOrdinal int + commitSnapshotID int64 +} + +// NewDeletedRowsScanTask constructs a row-level delete task. existingDeletes +// are stored on the embedded FileScanTask as DeleteFiles so later readers can +// reuse the normal scan delete path for the live-row baseline. +func NewDeletedRowsScanTask(dataFile iceberg.DataFile, addedDeletes, existingDeletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) DeletedRowsScanTask { + return DeletedRowsScanTask{ + FileScanTask: FileScanTask{ + File: dataFile, + DeleteFiles: existingDeletes, + }, + addedDeletes: addedDeletes, + changeOrdinal: changeOrdinal, + commitSnapshotID: commitSnapshotID, + } +} + +func (t DeletedRowsScanTask) ChangeOrdinal() int { return t.changeOrdinal } +func (t DeletedRowsScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } + +// AddedDeletes returns delete files whose removals should appear in the +// changelog. +func (t DeletedRowsScanTask) AddedDeletes() []iceberg.DataFile { + return t.addedDeletes +} + +// ExistingDeletes returns delete files that already applied before this +// snapshot's added deletes. +func (t DeletedRowsScanTask) ExistingDeletes() []iceberg.DataFile { + return t.DeleteFiles +} diff --git a/table/changelog_scan_task_test.go b/table/changelog_scan_task_test.go new file mode 100644 index 000000000..5995f2fc6 --- /dev/null +++ b/table/changelog_scan_task_test.go @@ -0,0 +1,72 @@ +// 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 table + +import ( + "testing" + + "github.com/apache/iceberg-go" + "github.com/stretchr/testify/require" +) + +func changelogTestDataFile(t *testing.T, path string, content iceberg.ManifestEntryContent) iceberg.DataFile { + t.Helper() + + b, err := iceberg.NewDataFileBuilder(*iceberg.UnpartitionedSpec, + content, path, iceberg.ParquetFile, nil, nil, nil, 10, 1024) + require.NoError(t, err) + + return b.Build() +} + +func TestAddedRowsScanTaskAppliesSameSnapshotDeletes(t *testing.T) { + data := changelogTestDataFile(t, "data/f1.parquet", iceberg.EntryContentData) + posDel := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes) + + task := NewAddedRowsScanTask(data, []iceberg.DataFile{posDel}, 0, 42) + + require.Equal(t, 0, task.ChangeOrdinal()) + require.Equal(t, int64(42), task.CommitSnapshotID()) + require.Equal(t, data.FilePath(), task.File.FilePath()) + require.Equal(t, []iceberg.DataFile{posDel}, task.Deletes()) +} + +func TestDeletedDataFileScanTaskKeepsExistingDeletes(t *testing.T) { + data := changelogTestDataFile(t, "data/f2.parquet", iceberg.EntryContentData) + existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes) + + task := NewDeletedDataFileScanTask(data, []iceberg.DataFile{existing}, 1, 43) + + require.Equal(t, 1, task.ChangeOrdinal()) + require.Equal(t, int64(43), task.CommitSnapshotID()) + require.Equal(t, []iceberg.DataFile{existing}, task.ExistingDeletes()) +} + +func TestDeletedRowsScanTaskSeparatesAddedAndExistingDeletes(t *testing.T) { + data := changelogTestDataFile(t, "data/f2.parquet", iceberg.EntryContentData) + added := changelogTestDataFile(t, "deletes/d2.parquet", iceberg.EntryContentEqDeletes) + existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes) + + task := NewDeletedRowsScanTask(data, []iceberg.DataFile{added}, []iceberg.DataFile{existing}, 2, 44) + + require.Equal(t, 2, task.ChangeOrdinal()) + require.Equal(t, int64(44), task.CommitSnapshotID()) + require.Equal(t, []iceberg.DataFile{added}, task.AddedDeletes()) + require.Equal(t, []iceberg.DataFile{existing}, task.ExistingDeletes()) + require.Equal(t, existing.FilePath(), task.DeleteFiles[0].FilePath()) +} From d20d3661dfbd78c6f890430d86afffe7a54fd0d1 Mon Sep 17 00:00:00 2001 From: Digvijay Date: Thu, 27 Aug 2026 00:05:52 -0500 Subject: [PATCH 2/3] fix(table): split changelog deletes by kind Store position deletes, equality deletes, and deletion vectors on the matching FileScanTask fields instead of collapsing every delete into positional DeleteFiles. Signed-off-by: Digvijay --- table/changelog_scan_task.go | 73 +++++++++++++++++++++---------- table/changelog_scan_task_test.go | 30 ++++++++----- 2 files changed, 70 insertions(+), 33 deletions(-) diff --git a/table/changelog_scan_task.go b/table/changelog_scan_task.go index 279c368e6..395e45b52 100644 --- a/table/changelog_scan_task.go +++ b/table/changelog_scan_task.go @@ -30,13 +30,12 @@ type AddedRowsScanTask struct { } // NewAddedRowsScanTask constructs an insert task for dataFile. deletes are -// delete files that apply while reading the added file. +// delete files that apply while reading the added file. Position deletes, +// equality deletes, and deletion vectors are stored on the matching +// FileScanTask fields. func NewAddedRowsScanTask(dataFile iceberg.DataFile, deletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) AddedRowsScanTask { return AddedRowsScanTask{ - FileScanTask: FileScanTask{ - File: dataFile, - DeleteFiles: deletes, - }, + FileScanTask: fileScanTaskWithDeletes(dataFile, deletes), changeOrdinal: changeOrdinal, commitSnapshotID: commitSnapshotID, } @@ -45,9 +44,10 @@ func NewAddedRowsScanTask(dataFile iceberg.DataFile, deletes []iceberg.DataFile, func (t AddedRowsScanTask) ChangeOrdinal() int { return t.changeOrdinal } func (t AddedRowsScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } -// Deletes returns delete files to apply when reading the added data file. +// Deletes returns every delete file applied while reading the added data +// file: position deletes, then equality deletes, then deletion vectors. func (t AddedRowsScanTask) Deletes() []iceberg.DataFile { - return t.DeleteFiles + return allDeleteFiles(t.FileScanTask) } // DeletedDataFileScanTask is a changelog delete produced by removing a data @@ -63,10 +63,7 @@ type DeletedDataFileScanTask struct { // NewDeletedDataFileScanTask constructs a delete task for a removed data file. func NewDeletedDataFileScanTask(dataFile iceberg.DataFile, existingDeletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) DeletedDataFileScanTask { return DeletedDataFileScanTask{ - FileScanTask: FileScanTask{ - File: dataFile, - DeleteFiles: existingDeletes, - }, + FileScanTask: fileScanTaskWithDeletes(dataFile, existingDeletes), changeOrdinal: changeOrdinal, commitSnapshotID: commitSnapshotID, } @@ -78,7 +75,7 @@ func (t DeletedDataFileScanTask) CommitSnapshotID() int64 { return t.commitSnaps // ExistingDeletes returns delete files that applied before the data file was // removed. func (t DeletedDataFileScanTask) ExistingDeletes() []iceberg.DataFile { - return t.DeleteFiles + return allDeleteFiles(t.FileScanTask) } // DeletedRowsScanTask is a changelog delete produced by adding delete files @@ -87,21 +84,18 @@ func (t DeletedDataFileScanTask) ExistingDeletes() []iceberg.DataFile { // those rows must not be emitted again. type DeletedRowsScanTask struct { FileScanTask - addedDeletes []iceberg.DataFile + addedDeletes FileScanTask changeOrdinal int commitSnapshotID int64 } // NewDeletedRowsScanTask constructs a row-level delete task. existingDeletes -// are stored on the embedded FileScanTask as DeleteFiles so later readers can -// reuse the normal scan delete path for the live-row baseline. +// are stored on the embedded FileScanTask so later readers can reuse the +// normal scan delete path for the live-row baseline. func NewDeletedRowsScanTask(dataFile iceberg.DataFile, addedDeletes, existingDeletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) DeletedRowsScanTask { return DeletedRowsScanTask{ - FileScanTask: FileScanTask{ - File: dataFile, - DeleteFiles: existingDeletes, - }, - addedDeletes: addedDeletes, + FileScanTask: fileScanTaskWithDeletes(dataFile, existingDeletes), + addedDeletes: fileScanTaskWithDeletes(dataFile, addedDeletes), changeOrdinal: changeOrdinal, commitSnapshotID: commitSnapshotID, } @@ -113,11 +107,46 @@ func (t DeletedRowsScanTask) CommitSnapshotID() int64 { return t.commitSnapshotI // AddedDeletes returns delete files whose removals should appear in the // changelog. func (t DeletedRowsScanTask) AddedDeletes() []iceberg.DataFile { - return t.addedDeletes + return allDeleteFiles(t.addedDeletes) } // ExistingDeletes returns delete files that already applied before this // snapshot's added deletes. func (t DeletedRowsScanTask) ExistingDeletes() []iceberg.DataFile { - return t.DeleteFiles + return allDeleteFiles(t.FileScanTask) +} + +func fileScanTaskWithDeletes(dataFile iceberg.DataFile, deletes []iceberg.DataFile) FileScanTask { + pos, eq, dv := classifyDeleteFiles(deletes) + return FileScanTask{ + File: dataFile, + DeleteFiles: pos, + EqualityDeleteFiles: eq, + DeletionVectorFiles: dv, + } +} + +func classifyDeleteFiles(files []iceberg.DataFile) (pos, eq, dv []iceberg.DataFile) { + for _, f := range files { + if f == nil { + continue + } + switch { + case IsDeletionVector(f): + dv = append(dv, f) + case f.ContentType() == iceberg.EntryContentEqDeletes: + eq = append(eq, f) + case f.ContentType() == iceberg.EntryContentPosDeletes: + pos = append(pos, f) + } + } + return pos, eq, dv +} + +func allDeleteFiles(task FileScanTask) []iceberg.DataFile { + out := make([]iceberg.DataFile, 0, len(task.DeleteFiles)+len(task.EqualityDeleteFiles)+len(task.DeletionVectorFiles)) + out = append(out, task.DeleteFiles...) + out = append(out, task.EqualityDeleteFiles...) + out = append(out, task.DeletionVectorFiles...) + return out } diff --git a/table/changelog_scan_task_test.go b/table/changelog_scan_task_test.go index 5995f2fc6..911e6e217 100644 --- a/table/changelog_scan_task_test.go +++ b/table/changelog_scan_task_test.go @@ -24,43 +24,49 @@ import ( "github.com/stretchr/testify/require" ) -func changelogTestDataFile(t *testing.T, path string, content iceberg.ManifestEntryContent) iceberg.DataFile { +func changelogTestDataFile(t *testing.T, path string, content iceberg.ManifestEntryContent, format iceberg.FileFormat) iceberg.DataFile { t.Helper() b, err := iceberg.NewDataFileBuilder(*iceberg.UnpartitionedSpec, - content, path, iceberg.ParquetFile, nil, nil, nil, 10, 1024) + content, path, format, nil, nil, nil, 10, 1024) require.NoError(t, err) return b.Build() } func TestAddedRowsScanTaskAppliesSameSnapshotDeletes(t *testing.T) { - data := changelogTestDataFile(t, "data/f1.parquet", iceberg.EntryContentData) - posDel := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes) + data := changelogTestDataFile(t, "data/f1.parquet", iceberg.EntryContentData, iceberg.ParquetFile) + posDel := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes, iceberg.ParquetFile) + eqDel := changelogTestDataFile(t, "deletes/d2.parquet", iceberg.EntryContentEqDeletes, iceberg.ParquetFile) + dv := changelogTestDataFile(t, "deletes/d3.puffin", iceberg.EntryContentPosDeletes, iceberg.PuffinFile) - task := NewAddedRowsScanTask(data, []iceberg.DataFile{posDel}, 0, 42) + task := NewAddedRowsScanTask(data, []iceberg.DataFile{eqDel, dv, posDel}, 0, 42) require.Equal(t, 0, task.ChangeOrdinal()) require.Equal(t, int64(42), task.CommitSnapshotID()) require.Equal(t, data.FilePath(), task.File.FilePath()) - require.Equal(t, []iceberg.DataFile{posDel}, task.Deletes()) + require.Equal(t, []iceberg.DataFile{posDel}, task.DeleteFiles) + require.Equal(t, []iceberg.DataFile{eqDel}, task.EqualityDeleteFiles) + require.Equal(t, []iceberg.DataFile{dv}, task.DeletionVectorFiles) + require.Equal(t, []iceberg.DataFile{posDel, eqDel, dv}, task.Deletes()) } func TestDeletedDataFileScanTaskKeepsExistingDeletes(t *testing.T) { - data := changelogTestDataFile(t, "data/f2.parquet", iceberg.EntryContentData) - existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes) + data := changelogTestDataFile(t, "data/f2.parquet", iceberg.EntryContentData, iceberg.ParquetFile) + existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes, iceberg.ParquetFile) task := NewDeletedDataFileScanTask(data, []iceberg.DataFile{existing}, 1, 43) require.Equal(t, 1, task.ChangeOrdinal()) require.Equal(t, int64(43), task.CommitSnapshotID()) require.Equal(t, []iceberg.DataFile{existing}, task.ExistingDeletes()) + require.Equal(t, []iceberg.DataFile{existing}, task.DeleteFiles) } func TestDeletedRowsScanTaskSeparatesAddedAndExistingDeletes(t *testing.T) { - data := changelogTestDataFile(t, "data/f2.parquet", iceberg.EntryContentData) - added := changelogTestDataFile(t, "deletes/d2.parquet", iceberg.EntryContentEqDeletes) - existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes) + data := changelogTestDataFile(t, "data/f2.parquet", iceberg.EntryContentData, iceberg.ParquetFile) + added := changelogTestDataFile(t, "deletes/d2.parquet", iceberg.EntryContentEqDeletes, iceberg.ParquetFile) + existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes, iceberg.ParquetFile) task := NewDeletedRowsScanTask(data, []iceberg.DataFile{added}, []iceberg.DataFile{existing}, 2, 44) @@ -69,4 +75,6 @@ func TestDeletedRowsScanTaskSeparatesAddedAndExistingDeletes(t *testing.T) { require.Equal(t, []iceberg.DataFile{added}, task.AddedDeletes()) require.Equal(t, []iceberg.DataFile{existing}, task.ExistingDeletes()) require.Equal(t, existing.FilePath(), task.DeleteFiles[0].FilePath()) + require.Empty(t, task.EqualityDeleteFiles) + require.Equal(t, []iceberg.DataFile{added}, task.addedDeletes.EqualityDeleteFiles) } From 89c1c77df14685adcb69aa57ed43f7245c0cecdb Mon Sep 17 00:00:00 2001 From: Digvijay Date: Thu, 27 Aug 2026 17:08:09 -0500 Subject: [PATCH 3/3] fix(table): align changelog tasks with review feedback Add ChangelogScanTask and Operation(), store addedDeletes as classified lists instead of a FileScanTask, and error on unknown delete content. Signed-off-by: Digvijay --- table/changelog_scan_task.go | 158 ++++++++++++++++++++++-------- table/changelog_scan_task_test.go | 46 ++++++++- table/scanner.go | 54 +++++++--- 3 files changed, 200 insertions(+), 58 deletions(-) diff --git a/table/changelog_scan_task.go b/table/changelog_scan_task.go index 395e45b52..bfa5a4106 100644 --- a/table/changelog_scan_task.go +++ b/table/changelog_scan_task.go @@ -17,7 +17,49 @@ package table -import "github.com/apache/iceberg-go" +import ( + "fmt" + + "github.com/apache/iceberg-go" +) + +// ChangelogOperation is the kind of change a changelog scan task produces. +type ChangelogOperation string + +const ( + ChangelogOpInsert ChangelogOperation = "INSERT" + ChangelogOpDelete ChangelogOperation = "DELETE" + ChangelogOpUpdateBefore ChangelogOperation = "UPDATE_BEFORE" + ChangelogOpUpdateAfter ChangelogOperation = "UPDATE_AFTER" +) + +// ChangelogScanTask is a unit of work that produces changelog rows. +type ChangelogScanTask interface { + Operation() ChangelogOperation + ChangeOrdinal() int + CommitSnapshotID() int64 +} + +var ( + _ ChangelogScanTask = AddedRowsScanTask{} + _ ChangelogScanTask = DeletedDataFileScanTask{} + _ ChangelogScanTask = DeletedRowsScanTask{} +) + +// classifiedDeletes holds delete files split the same way FileScanTask does, +// without a second FileScanTask whose range and lineage fields would be zero. +type classifiedDeletes struct { + pos, eq, dv []iceberg.DataFile +} + +func (d classifiedDeletes) files() []iceberg.DataFile { + out := make([]iceberg.DataFile, 0, len(d.pos)+len(d.eq)+len(d.dv)) + out = append(out, d.pos...) + out = append(out, d.eq...) + out = append(out, d.dv...) + + return out +} // AddedRowsScanTask is a changelog insert produced by adding a data file. // Matching delete files committed in the same snapshot, or from squashed @@ -33,16 +75,22 @@ type AddedRowsScanTask struct { // delete files that apply while reading the added file. Position deletes, // equality deletes, and deletion vectors are stored on the matching // FileScanTask fields. -func NewAddedRowsScanTask(dataFile iceberg.DataFile, deletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) AddedRowsScanTask { +func NewAddedRowsScanTask(dataFile iceberg.DataFile, deletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) (AddedRowsScanTask, error) { + task, err := fileScanTaskWithDeletes(dataFile, deletes) + if err != nil { + return AddedRowsScanTask{}, err + } + return AddedRowsScanTask{ - FileScanTask: fileScanTaskWithDeletes(dataFile, deletes), + FileScanTask: task, changeOrdinal: changeOrdinal, commitSnapshotID: commitSnapshotID, - } + }, nil } -func (t AddedRowsScanTask) ChangeOrdinal() int { return t.changeOrdinal } -func (t AddedRowsScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } +func (t AddedRowsScanTask) Operation() ChangelogOperation { return ChangelogOpInsert } +func (t AddedRowsScanTask) ChangeOrdinal() int { return t.changeOrdinal } +func (t AddedRowsScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } // Deletes returns every delete file applied while reading the added data // file: position deletes, then equality deletes, then deletion vectors. @@ -61,16 +109,22 @@ type DeletedDataFileScanTask struct { } // NewDeletedDataFileScanTask constructs a delete task for a removed data file. -func NewDeletedDataFileScanTask(dataFile iceberg.DataFile, existingDeletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) DeletedDataFileScanTask { +func NewDeletedDataFileScanTask(dataFile iceberg.DataFile, existingDeletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) (DeletedDataFileScanTask, error) { + task, err := fileScanTaskWithDeletes(dataFile, existingDeletes) + if err != nil { + return DeletedDataFileScanTask{}, err + } + return DeletedDataFileScanTask{ - FileScanTask: fileScanTaskWithDeletes(dataFile, existingDeletes), + FileScanTask: task, changeOrdinal: changeOrdinal, commitSnapshotID: commitSnapshotID, - } + }, nil } -func (t DeletedDataFileScanTask) ChangeOrdinal() int { return t.changeOrdinal } -func (t DeletedDataFileScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } +func (t DeletedDataFileScanTask) Operation() ChangelogOperation { return ChangelogOpDelete } +func (t DeletedDataFileScanTask) ChangeOrdinal() int { return t.changeOrdinal } +func (t DeletedDataFileScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } // ExistingDeletes returns delete files that applied before the data file was // removed. @@ -84,7 +138,7 @@ func (t DeletedDataFileScanTask) ExistingDeletes() []iceberg.DataFile { // those rows must not be emitted again. type DeletedRowsScanTask struct { FileScanTask - addedDeletes FileScanTask + addedDeletes classifiedDeletes changeOrdinal int commitSnapshotID int64 } @@ -92,22 +146,33 @@ type DeletedRowsScanTask struct { // NewDeletedRowsScanTask constructs a row-level delete task. existingDeletes // are stored on the embedded FileScanTask so later readers can reuse the // normal scan delete path for the live-row baseline. -func NewDeletedRowsScanTask(dataFile iceberg.DataFile, addedDeletes, existingDeletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) DeletedRowsScanTask { +func NewDeletedRowsScanTask(dataFile iceberg.DataFile, addedDeletes, existingDeletes []iceberg.DataFile, changeOrdinal int, commitSnapshotID int64) (DeletedRowsScanTask, error) { + existing, err := fileScanTaskWithDeletes(dataFile, existingDeletes) + if err != nil { + return DeletedRowsScanTask{}, err + } + + added, err := classifyDeleteFiles(addedDeletes) + if err != nil { + return DeletedRowsScanTask{}, err + } + return DeletedRowsScanTask{ - FileScanTask: fileScanTaskWithDeletes(dataFile, existingDeletes), - addedDeletes: fileScanTaskWithDeletes(dataFile, addedDeletes), + FileScanTask: existing, + addedDeletes: added, changeOrdinal: changeOrdinal, commitSnapshotID: commitSnapshotID, - } + }, nil } -func (t DeletedRowsScanTask) ChangeOrdinal() int { return t.changeOrdinal } -func (t DeletedRowsScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } +func (t DeletedRowsScanTask) Operation() ChangelogOperation { return ChangelogOpDelete } +func (t DeletedRowsScanTask) ChangeOrdinal() int { return t.changeOrdinal } +func (t DeletedRowsScanTask) CommitSnapshotID() int64 { return t.commitSnapshotID } // AddedDeletes returns delete files whose removals should appear in the // changelog. func (t DeletedRowsScanTask) AddedDeletes() []iceberg.DataFile { - return allDeleteFiles(t.addedDeletes) + return t.addedDeletes.files() } // ExistingDeletes returns delete files that already applied before this @@ -116,37 +181,48 @@ func (t DeletedRowsScanTask) ExistingDeletes() []iceberg.DataFile { return allDeleteFiles(t.FileScanTask) } -func fileScanTaskWithDeletes(dataFile iceberg.DataFile, deletes []iceberg.DataFile) FileScanTask { - pos, eq, dv := classifyDeleteFiles(deletes) +func fileScanTaskWithDeletes(dataFile iceberg.DataFile, deletes []iceberg.DataFile) (FileScanTask, error) { + classified, err := classifyDeleteFiles(deletes) + if err != nil { + return FileScanTask{}, err + } + return FileScanTask{ File: dataFile, - DeleteFiles: pos, - EqualityDeleteFiles: eq, - DeletionVectorFiles: dv, - } + DeleteFiles: classified.pos, + EqualityDeleteFiles: classified.eq, + DeletionVectorFiles: classified.dv, + }, nil } -func classifyDeleteFiles(files []iceberg.DataFile) (pos, eq, dv []iceberg.DataFile) { +func classifyDeleteFiles(files []iceberg.DataFile) (classifiedDeletes, error) { + var out classifiedDeletes for _, f := range files { - if f == nil { - continue + kind, err := classifyDataFile(f) + if err != nil { + return classifiedDeletes{}, err } - switch { - case IsDeletionVector(f): - dv = append(dv, f) - case f.ContentType() == iceberg.EntryContentEqDeletes: - eq = append(eq, f) - case f.ContentType() == iceberg.EntryContentPosDeletes: - pos = append(pos, f) + + switch kind { + case dataFileKindPosDeletes: + out.pos = append(out.pos, f) + case dataFileKindEqDeletes: + out.eq = append(out.eq, f) + case dataFileKindDeletionVector: + out.dv = append(out.dv, f) + default: + return classifiedDeletes{}, fmt.Errorf("%w: expected delete file, got content type %s", + ErrInvalidMetadata, f.ContentType()) } } - return pos, eq, dv + + return out, nil } func allDeleteFiles(task FileScanTask) []iceberg.DataFile { - out := make([]iceberg.DataFile, 0, len(task.DeleteFiles)+len(task.EqualityDeleteFiles)+len(task.DeletionVectorFiles)) - out = append(out, task.DeleteFiles...) - out = append(out, task.EqualityDeleteFiles...) - out = append(out, task.DeletionVectorFiles...) - return out + return classifiedDeletes{ + pos: task.DeleteFiles, + eq: task.EqualityDeleteFiles, + dv: task.DeletionVectorFiles, + }.files() } diff --git a/table/changelog_scan_task_test.go b/table/changelog_scan_task_test.go index 911e6e217..e5af7f5ab 100644 --- a/table/changelog_scan_task_test.go +++ b/table/changelog_scan_task_test.go @@ -40,8 +40,10 @@ func TestAddedRowsScanTaskAppliesSameSnapshotDeletes(t *testing.T) { eqDel := changelogTestDataFile(t, "deletes/d2.parquet", iceberg.EntryContentEqDeletes, iceberg.ParquetFile) dv := changelogTestDataFile(t, "deletes/d3.puffin", iceberg.EntryContentPosDeletes, iceberg.PuffinFile) - task := NewAddedRowsScanTask(data, []iceberg.DataFile{eqDel, dv, posDel}, 0, 42) + task, err := NewAddedRowsScanTask(data, []iceberg.DataFile{eqDel, dv, posDel}, 0, 42) + require.NoError(t, err) + require.Equal(t, ChangelogOpInsert, task.Operation()) require.Equal(t, 0, task.ChangeOrdinal()) require.Equal(t, int64(42), task.CommitSnapshotID()) require.Equal(t, data.FilePath(), task.File.FilePath()) @@ -55,8 +57,10 @@ func TestDeletedDataFileScanTaskKeepsExistingDeletes(t *testing.T) { data := changelogTestDataFile(t, "data/f2.parquet", iceberg.EntryContentData, iceberg.ParquetFile) existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes, iceberg.ParquetFile) - task := NewDeletedDataFileScanTask(data, []iceberg.DataFile{existing}, 1, 43) + task, err := NewDeletedDataFileScanTask(data, []iceberg.DataFile{existing}, 1, 43) + require.NoError(t, err) + require.Equal(t, ChangelogOpDelete, task.Operation()) require.Equal(t, 1, task.ChangeOrdinal()) require.Equal(t, int64(43), task.CommitSnapshotID()) require.Equal(t, []iceberg.DataFile{existing}, task.ExistingDeletes()) @@ -68,13 +72,47 @@ func TestDeletedRowsScanTaskSeparatesAddedAndExistingDeletes(t *testing.T) { added := changelogTestDataFile(t, "deletes/d2.parquet", iceberg.EntryContentEqDeletes, iceberg.ParquetFile) existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes, iceberg.ParquetFile) - task := NewDeletedRowsScanTask(data, []iceberg.DataFile{added}, []iceberg.DataFile{existing}, 2, 44) + task, err := NewDeletedRowsScanTask(data, []iceberg.DataFile{added}, []iceberg.DataFile{existing}, 2, 44) + require.NoError(t, err) + require.Equal(t, ChangelogOpDelete, task.Operation()) require.Equal(t, 2, task.ChangeOrdinal()) require.Equal(t, int64(44), task.CommitSnapshotID()) require.Equal(t, []iceberg.DataFile{added}, task.AddedDeletes()) require.Equal(t, []iceberg.DataFile{existing}, task.ExistingDeletes()) require.Equal(t, existing.FilePath(), task.DeleteFiles[0].FilePath()) require.Empty(t, task.EqualityDeleteFiles) - require.Equal(t, []iceberg.DataFile{added}, task.addedDeletes.EqualityDeleteFiles) +} + +func TestChangelogScanTaskInterface(t *testing.T) { + data := changelogTestDataFile(t, "data/f1.parquet", iceberg.EntryContentData, iceberg.ParquetFile) + + added, err := NewAddedRowsScanTask(data, nil, 0, 1) + require.NoError(t, err) + deletedFile, err := NewDeletedDataFileScanTask(data, nil, 1, 2) + require.NoError(t, err) + deletedRows, err := NewDeletedRowsScanTask(data, nil, nil, 2, 3) + require.NoError(t, err) + + tasks := []ChangelogScanTask{added, deletedFile, deletedRows} + require.Equal(t, ChangelogOpInsert, tasks[0].Operation()) + require.Equal(t, ChangelogOpDelete, tasks[1].Operation()) + require.Equal(t, ChangelogOpDelete, tasks[2].Operation()) +} + +func TestClassifyDeleteFiles(t *testing.T) { + posDel := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes, iceberg.ParquetFile) + eqDel := changelogTestDataFile(t, "deletes/d2.parquet", iceberg.EntryContentEqDeletes, iceberg.ParquetFile) + dv := changelogTestDataFile(t, "deletes/d3.puffin", iceberg.EntryContentPosDeletes, iceberg.PuffinFile) + data := changelogTestDataFile(t, "data/f1.parquet", iceberg.EntryContentData, iceberg.ParquetFile) + + got, err := classifyDeleteFiles([]iceberg.DataFile{eqDel, dv, posDel}) + require.NoError(t, err) + require.Equal(t, []iceberg.DataFile{posDel}, got.pos) + require.Equal(t, []iceberg.DataFile{eqDel}, got.eq) + require.Equal(t, []iceberg.DataFile{dv}, got.dv) + + _, err = classifyDeleteFiles([]iceberg.DataFile{data}) + require.ErrorIs(t, err, ErrInvalidMetadata) + require.ErrorContains(t, err, "expected delete file") } diff --git a/table/scanner.go b/table/scanner.go index a66601b4f..ace8f2c4c 100644 --- a/table/scanner.go +++ b/table/scanner.go @@ -151,21 +151,20 @@ func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error { defer m.mu.Unlock() for _, entry := range entries { - dataFile := entry.DataFile() - switch dataFile.ContentType() { - case iceberg.EntryContentData: + kind, err := classifyDataFile(entry.DataFile()) + if err != nil { + return fmt.Errorf("%w: %s", err, entry) + } + + switch kind { + case dataFileKindData: m.dataEntries = append(m.dataEntries, entry) - case iceberg.EntryContentPosDeletes: - if IsDeletionVector(dataFile) { - m.dvEntries = append(m.dvEntries, entry) - } else { - m.positionalDeleteEntries = append(m.positionalDeleteEntries, entry) - } - case iceberg.EntryContentEqDeletes: + case dataFileKindPosDeletes: + m.positionalDeleteEntries = append(m.positionalDeleteEntries, entry) + case dataFileKindEqDeletes: m.equalityDeleteEntries = append(m.equalityDeleteEntries, entry) - default: - return fmt.Errorf("%w: unknown DataFileContent type (%s): %s", - ErrInvalidMetadata, dataFile.ContentType(), entry) + case dataFileKindDeletionVector: + m.dvEntries = append(m.dvEntries, entry) } } @@ -230,6 +229,35 @@ func IsDeletionVector(df iceberg.DataFile) bool { df.ContentType() == iceberg.EntryContentPosDeletes } +type dataFileKind int + +const ( + dataFileKindData dataFileKind = iota + dataFileKindPosDeletes + dataFileKindEqDeletes + dataFileKindDeletionVector +) + +// classifyDataFile buckets a file by content type. Deletion vectors are +// Puffin position-delete files and are split out from regular pos-deletes. +func classifyDataFile(f iceberg.DataFile) (dataFileKind, error) { + switch f.ContentType() { + case iceberg.EntryContentData: + return dataFileKindData, nil + case iceberg.EntryContentPosDeletes: + if IsDeletionVector(f) { + return dataFileKindDeletionVector, nil + } + + return dataFileKindPosDeletes, nil + case iceberg.EntryContentEqDeletes: + return dataFileKindEqDeletes, nil + default: + return 0, fmt.Errorf("%w: unknown DataFileContent type (%s)", + ErrInvalidMetadata, f.ContentType()) + } +} + type Scan struct { identifier Identifier metadata Metadata