diff --git a/table/changelog_scan_task.go b/table/changelog_scan_task.go new file mode 100644 index 000000000..bfa5a4106 --- /dev/null +++ b/table/changelog_scan_task.go @@ -0,0 +1,228 @@ +// 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 ( + "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 +// 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. 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, error) { + task, err := fileScanTaskWithDeletes(dataFile, deletes) + if err != nil { + return AddedRowsScanTask{}, err + } + + return AddedRowsScanTask{ + FileScanTask: task, + changeOrdinal: changeOrdinal, + commitSnapshotID: commitSnapshotID, + }, nil +} + +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. +func (t AddedRowsScanTask) Deletes() []iceberg.DataFile { + return allDeleteFiles(t.FileScanTask) +} + +// 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, error) { + task, err := fileScanTaskWithDeletes(dataFile, existingDeletes) + if err != nil { + return DeletedDataFileScanTask{}, err + } + + return DeletedDataFileScanTask{ + FileScanTask: task, + changeOrdinal: changeOrdinal, + commitSnapshotID: commitSnapshotID, + }, nil +} + +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. +func (t DeletedDataFileScanTask) ExistingDeletes() []iceberg.DataFile { + return allDeleteFiles(t.FileScanTask) +} + +// 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 classifiedDeletes + changeOrdinal int + commitSnapshotID int64 +} + +// 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, 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: existing, + addedDeletes: added, + changeOrdinal: changeOrdinal, + commitSnapshotID: commitSnapshotID, + }, nil +} + +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 t.addedDeletes.files() +} + +// ExistingDeletes returns delete files that already applied before this +// snapshot's added deletes. +func (t DeletedRowsScanTask) ExistingDeletes() []iceberg.DataFile { + return allDeleteFiles(t.FileScanTask) +} + +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: classified.pos, + EqualityDeleteFiles: classified.eq, + DeletionVectorFiles: classified.dv, + }, nil +} + +func classifyDeleteFiles(files []iceberg.DataFile) (classifiedDeletes, error) { + var out classifiedDeletes + for _, f := range files { + kind, err := classifyDataFile(f) + if err != nil { + return classifiedDeletes{}, err + } + + 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 out, nil +} + +func allDeleteFiles(task FileScanTask) []iceberg.DataFile { + 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 new file mode 100644 index 000000000..e5af7f5ab --- /dev/null +++ b/table/changelog_scan_task_test.go @@ -0,0 +1,118 @@ +// 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, format iceberg.FileFormat) iceberg.DataFile { + t.Helper() + + b, err := iceberg.NewDataFileBuilder(*iceberg.UnpartitionedSpec, + 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, 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, 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()) + 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, iceberg.ParquetFile) + existing := changelogTestDataFile(t, "deletes/d1.parquet", iceberg.EntryContentPosDeletes, iceberg.ParquetFile) + + 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()) + require.Equal(t, []iceberg.DataFile{existing}, task.DeleteFiles) +} + +func TestDeletedRowsScanTaskSeparatesAddedAndExistingDeletes(t *testing.T) { + 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, 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) +} + +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 f5599d932..4cbd6d2f0 100644 --- a/table/scanner.go +++ b/table/scanner.go @@ -376,6 +376,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()) + } +} + // Scan represents a table scan. It implements [io.Closer]; callers should // close it when they are done, including early exits after remote planning // succeeds but before all records are consumed.