laskoviymishka commented on code in PR #2046:
URL: https://github.com/apache/iceberg-go/pull/2046#discussion_r4075520452
##########
table/transaction.go:
##########
@@ -2456,7 +2455,7 @@ func (t *Transaction) performCopyOnWriteDeletion(ctx
context.Context, operation
updater := t.updateSnapshot(wfs, snapshotProps,
operation).mergeOverwrite(&commitUUID, filter)
updater.setManifestConcurrency(concurrency)
- filesToDelete, filesToRewrite, fileSeqByPath, err :=
t.classifyFilesForDeletions(ctx, fs, filter, caseSensitive, concurrency)
+ filesToDelete, filesToRewrite, err := t.classifyFilesForDeletions(ctx,
fs, filter, caseSensitive, concurrency)
Review Comment:
tiny thing: `classifyFilesForDeletions` now returns
`[]iceberg.ManifestEntry` for this second value, but `filesToRewrite` still
reads like it holds `DataFile`s. `entriesToRewrite` would match the `entries`
name you're already using downstream in
`rewriteFilesWithFilter`/`rewriteScanTasks`.
##########
table/transaction.go:
##########
@@ -2800,17 +2785,72 @@ func (t *Transaction)
classifyFilesForFilteredDeletions(ctx context.Context, fs
}
if err := g.Wait(); err != nil {
- return nil, nil, nil, err
+ return nil, nil, err
}
- return filesToDelete, filesWithPartialDeletes, fileSeqByPath, nil
+ return filesToDelete, filesWithPartialDeletes, nil
+}
+
+// rewriteScanTasks returns a scan task for each rewrite candidate that
carries the deletes applying to it on the
+// planning snapshot, matched the same way scan planning matches them.
+//
+// Delete manifests are not pruned by the row filter: with filter id=2 the
rewrite keeps id=1,
+// so a delete on id=1 must still be applied.
+func (t *Transaction) rewriteScanTasks(fs io.IO, entries
[]iceberg.ManifestEntry) ([]FileScanTask, error) {
+ meta, err := t.txnMeta()
+ if err != nil {
+ return nil, err
+ }
+ builtMeta, err := meta.Build()
+ if err != nil {
+ return nil, err
+ }
+
+ var liveDeletes []iceberg.ManifestEntry
+ if s := t.planningSnapshot(meta); s != nil {
+ for entry, err := range s.entries(fs,
iceberg.ManifestContentDeletes) {
Review Comment:
this reads every live entry from every delete manifest in the planning
snapshot on each rewrite call, with no partition pruning. The doc comment is
right that we can't prune by the row filter (with filter id=2 we still need the
id=1 delete), but partition pruning is a separate lever: we already know the
exact partitions of `entries` here, so we could skip delete manifests whose
partition summary can't overlap any candidate, the way `planFilesLocal` does.
As written it also reimplements the planner's walk-classify-build-index
pipeline from scratch, so it's missing both that pruning and the `minSeqNum`
optimization, and it's a second copy of exactly the logic whose drift caused
#2042. Could we extract "collect classified delete entries + build the three
indices from a manifest list" into one helper shared by `planFilesLocal` and
`rewriteScanTasks`?
For a table with many delete files across many partitions, a narrow
single-file CoW delete otherwise pays to decode every delete manifest every
time. If scoping it is more than you want to take on here, I'd at least leave
an explicit comment (and maybe a follow-up issue) calling out the table-wide
read as a known limitation. wdyt?
##########
table/cow_rewrite_deletes_test.go:
##########
@@ -0,0 +1,171 @@
+// 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_test
+
+import (
+ "testing"
+
+ "github.com/apache/iceberg-go"
+ "github.com/apache/iceberg-go/table"
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+)
+
+// TestCoWRewriteKeepsRowsDeleted rewrites a data file that already has a
delete on id=1.
+// Once the file is replaced that delete no longer reaches the copied rows, so
the rewrite itself must drop id=1.
+func TestCoWRewriteKeepsRowsDeleted(t *testing.T) {
+ deleteID1 := func(t *testing.T, tbl *table.Table) *table.Table {
+ t.Helper()
+
+ out, err := tbl.Delete(t.Context(),
iceberg.EqualTo(iceberg.Reference("id"), int64(1)), nil)
+ require.NoError(t, err)
+
+ return out
+ }
+ equalityDeleteID1 := func(t *testing.T, tbl *table.Table) *table.Table {
+ t.Helper()
+
+ return appendEqualityDelete(t, tbl, []int{1}, `[{"id": 1}]`)
+ }
+
+ deletes := []struct {
+ name string
+ formatVersion string
+ seed func(*testing.T, *table.Table) *table.Table
+ }{
+ {name: "position delete", formatVersion: "2", seed: deleteID1},
+ {name: "deletion vector", formatVersion: "3", seed: deleteID1},
+ {name: "equality delete on v2", formatVersion: "2", seed:
equalityDeleteID1},
+ {name: "equality delete on v3", formatVersion: "3", seed:
equalityDeleteID1},
+ }
+ operations := []struct {
+ name string
+ run func(*testing.T, *table.Table) *table.Table
+ wantIDs []int64
+ }{
+ {
+ name: "copy-on-write delete",
+ run: func(t *testing.T, tbl *table.Table)
*table.Table { return copyOnWriteDelete(t, tbl, 2) },
+ wantIDs: []int64{3, 4, 5},
+ },
+ {
+ name: "filtered overwrite",
+ run: func(t *testing.T, tbl *table.Table)
*table.Table { return filteredOverwrite(t, tbl, 2, 6) },
+ wantIDs: []int64{3, 4, 5, 6},
+ },
+ }
+
+ for _, del := range deletes {
+ for _, op := range operations {
+ t.Run(del.name+"/"+op.name, func(t *testing.T) {
+ tbl := newMergeOnReadTestTableVersion(t,
del.formatVersion)
Review Comment:
all the new cases run against `newMergeOnReadTestTableVersion`, which is
unpartitioned, so we only exercise path-keyed delete matching through the new
path. On a partitioned table `rewriteScanTasks` also has to match deletes by
partition, which is exactly the matching most likely to drift from the planner.
A single partitioned variant of this matrix would cover it. Not blocking, but
worth adding while the test file is fresh.
##########
table/cow_rewrite_deletes_test.go:
##########
@@ -0,0 +1,171 @@
+// 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_test
+
+import (
+ "testing"
+
+ "github.com/apache/iceberg-go"
+ "github.com/apache/iceberg-go/table"
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+)
+
+// TestCoWRewriteKeepsRowsDeleted rewrites a data file that already has a
delete on id=1.
+// Once the file is replaced that delete no longer reaches the copied rows, so
the rewrite itself must drop id=1.
+func TestCoWRewriteKeepsRowsDeleted(t *testing.T) {
+ deleteID1 := func(t *testing.T, tbl *table.Table) *table.Table {
+ t.Helper()
+
+ out, err := tbl.Delete(t.Context(),
iceberg.EqualTo(iceberg.Reference("id"), int64(1)), nil)
+ require.NoError(t, err)
+
+ return out
+ }
+ equalityDeleteID1 := func(t *testing.T, tbl *table.Table) *table.Table {
+ t.Helper()
+
+ return appendEqualityDelete(t, tbl, []int{1}, `[{"id": 1}]`)
+ }
+
+ deletes := []struct {
+ name string
+ formatVersion string
+ seed func(*testing.T, *table.Table) *table.Table
+ }{
+ {name: "position delete", formatVersion: "2", seed: deleteID1},
+ {name: "deletion vector", formatVersion: "3", seed: deleteID1},
+ {name: "equality delete on v2", formatVersion: "2", seed:
equalityDeleteID1},
+ {name: "equality delete on v3", formatVersion: "3", seed:
equalityDeleteID1},
+ }
+ operations := []struct {
+ name string
+ run func(*testing.T, *table.Table) *table.Table
+ wantIDs []int64
+ }{
+ {
+ name: "copy-on-write delete",
+ run: func(t *testing.T, tbl *table.Table)
*table.Table { return copyOnWriteDelete(t, tbl, 2) },
+ wantIDs: []int64{3, 4, 5},
+ },
+ {
+ name: "filtered overwrite",
+ run: func(t *testing.T, tbl *table.Table)
*table.Table { return filteredOverwrite(t, tbl, 2, 6) },
+ wantIDs: []int64{3, 4, 5, 6},
+ },
+ }
+
+ for _, del := range deletes {
+ for _, op := range operations {
+ t.Run(del.name+"/"+op.name, func(t *testing.T) {
+ tbl := newMergeOnReadTestTableVersion(t,
del.formatVersion)
+ tbl = appendRowsOnRef(t, tbl, table.MainBranch,
1, 2, 3, 4, 5)
+ tbl = del.seed(t, tbl)
+ require.Equal(t, []int64{2, 3, 4, 5},
idsInTable(t, tbl))
+
+ var rowIDsBefore map[int64]int64
+ if del.formatVersion == "3" {
+ rowIDsBefore = readRowIDsByID(t,
t.Context(), tbl)
+ require.Len(t, rowIDsBefore, 4)
+ }
+
+ tbl = op.run(t, tbl)
+
+ assert.Equal(t, op.wantIDs, idsInTable(t, tbl),
"id=1 was deleted before the rewrite and must stay deleted")
+ assert.Equal(t, int64(len(op.wantIDs)),
liveDataRecordCount(t, tbl), "the rewritten file must not carry the
already-deleted row")
+
+ if rowIDsBefore != nil {
+ rowIDsAfter := readRowIDsByID(t,
t.Context(), tbl)
+ for _, id := range []int64{3, 4, 5} {
+ assert.Equal(t,
rowIDsBefore[id], rowIDsAfter[id],
+ "id=%d must keep its
_row_id through the rewrite", id)
+ }
+ }
+ })
+ }
+ }
+}
+
+// TestCoWRewriteSkipsDeletesThatDoNotApply re-inserts id=1 after an equality
delete on it.
+// The delete has a lower sequence number than the new file, so rewriting that
file must keep the row,
+// while the untouched older file must keep its own id=1 hidden.
+func TestCoWRewriteSkipsDeletesThatDoNotApply(t *testing.T) {
+ tbl := newMergeOnReadTestTableVersion(t, "2")
+ tbl = appendRowsOnRef(t, tbl, table.MainBranch, 1, 2, 3)
+ tbl = appendEqualityDelete(t, tbl, []int{1}, `[{"id": 1}]`)
+ tbl = appendRowsOnRef(t, tbl, table.MainBranch, 1, 4, 5)
+ require.Equal(t, []int64{1, 2, 3, 4, 5}, idsInTable(t, tbl))
+
+ tbl = copyOnWriteDelete(t, tbl, 4)
+
+ assert.Equal(t, []int64{1, 2, 3, 5}, idsInTable(t, tbl))
+ assert.Equal(t, int64(5), liveDataRecordCount(t, tbl), "only the file
holding id=4 is rewritten; the other keeps all 3 of its rows")
+}
+
+// copyOnWriteDelete switches the table to copy-on-write in its own commit,
then deletes id.
+func copyOnWriteDelete(t *testing.T, tbl *table.Table, id int64) *table.Table {
+ t.Helper()
+
+ tx := tbl.NewTransaction()
+ require.NoError(t,
tx.SetProperties(iceberg.Properties{table.WriteDeleteModeKey:
table.WriteModeCopyOnWrite}))
+ tbl, err := tx.Commit(t.Context())
+ require.NoError(t, err)
+
+ tbl, err = tbl.Delete(t.Context(),
iceberg.EqualTo(iceberg.Reference("id"), id), nil)
+ require.NoError(t, err)
+
+ return tbl
+}
+
+// filteredOverwrite replaces the rows matching id with a single row newID.
+func filteredOverwrite(t *testing.T, tbl *table.Table, id, newID int64)
*table.Table {
+ t.Helper()
+
+ rdr := idRecordReader(t, tbl, newID)
+ defer rdr.Release()
+
+ tbl, err := tbl.Overwrite(t.Context(), rdr, nil,
+
table.WithOverwriteFilter(iceberg.EqualTo(iceberg.Reference("id"), id)))
+ require.NoError(t, err)
+
+ return tbl
+}
+
+// liveDataRecordCount sums record_count over the live data files of the
current snapshot,
+// including rows that deletes hide from a scan.
+func liveDataRecordCount(t *testing.T, tbl *table.Table) int64 {
+ t.Helper()
+
+ fs, err := tbl.FS(t.Context())
+ require.NoError(t, err)
+ manifests, err := tbl.CurrentSnapshot().Manifests(fs)
Review Comment:
`CurrentSnapshot()` can return nil and `Manifests` has a value receiver, so
on a nil snapshot this panics instead of failing with a readable assertion. I'd
guard it:
```go
snap := tbl.CurrentSnapshot()
require.NotNil(t, snap)
manifests, err := snap.Manifests(fs)
```
##########
table/cow_rewrite_deletes_test.go:
##########
@@ -0,0 +1,171 @@
+// 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_test
+
+import (
+ "testing"
+
+ "github.com/apache/iceberg-go"
+ "github.com/apache/iceberg-go/table"
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+)
+
+// TestCoWRewriteKeepsRowsDeleted rewrites a data file that already has a
delete on id=1.
+// Once the file is replaced that delete no longer reaches the copied rows, so
the rewrite itself must drop id=1.
+func TestCoWRewriteKeepsRowsDeleted(t *testing.T) {
+ deleteID1 := func(t *testing.T, tbl *table.Table) *table.Table {
+ t.Helper()
+
+ out, err := tbl.Delete(t.Context(),
iceberg.EqualTo(iceberg.Reference("id"), int64(1)), nil)
+ require.NoError(t, err)
+
+ return out
+ }
+ equalityDeleteID1 := func(t *testing.T, tbl *table.Table) *table.Table {
+ t.Helper()
+
+ return appendEqualityDelete(t, tbl, []int{1}, `[{"id": 1}]`)
+ }
+
+ deletes := []struct {
+ name string
+ formatVersion string
+ seed func(*testing.T, *table.Table) *table.Table
+ }{
+ {name: "position delete", formatVersion: "2", seed: deleteID1},
+ {name: "deletion vector", formatVersion: "3", seed: deleteID1},
+ {name: "equality delete on v2", formatVersion: "2", seed:
equalityDeleteID1},
+ {name: "equality delete on v3", formatVersion: "3", seed:
equalityDeleteID1},
+ }
+ operations := []struct {
+ name string
+ run func(*testing.T, *table.Table) *table.Table
+ wantIDs []int64
+ }{
+ {
+ name: "copy-on-write delete",
+ run: func(t *testing.T, tbl *table.Table)
*table.Table { return copyOnWriteDelete(t, tbl, 2) },
+ wantIDs: []int64{3, 4, 5},
+ },
+ {
+ name: "filtered overwrite",
+ run: func(t *testing.T, tbl *table.Table)
*table.Table { return filteredOverwrite(t, tbl, 2, 6) },
+ wantIDs: []int64{3, 4, 5, 6},
+ },
+ }
+
+ for _, del := range deletes {
+ for _, op := range operations {
+ t.Run(del.name+"/"+op.name, func(t *testing.T) {
+ tbl := newMergeOnReadTestTableVersion(t,
del.formatVersion)
+ tbl = appendRowsOnRef(t, tbl, table.MainBranch,
1, 2, 3, 4, 5)
+ tbl = del.seed(t, tbl)
+ require.Equal(t, []int64{2, 3, 4, 5},
idsInTable(t, tbl))
+
+ var rowIDsBefore map[int64]int64
+ if del.formatVersion == "3" {
+ rowIDsBefore = readRowIDsByID(t,
t.Context(), tbl)
+ require.Len(t, rowIDsBefore, 4)
+ }
+
+ tbl = op.run(t, tbl)
+
+ assert.Equal(t, op.wantIDs, idsInTable(t, tbl),
"id=1 was deleted before the rewrite and must stay deleted")
+ assert.Equal(t, int64(len(op.wantIDs)),
liveDataRecordCount(t, tbl), "the rewritten file must not carry the
already-deleted row")
+
+ if rowIDsBefore != nil {
+ rowIDsAfter := readRowIDsByID(t,
t.Context(), tbl)
+ for _, id := range []int64{3, 4, 5} {
+ assert.Equal(t,
rowIDsBefore[id], rowIDsAfter[id],
+ "id=%d must keep its
_row_id through the rewrite", id)
+ }
+ }
+ })
+ }
+ }
+}
+
+// TestCoWRewriteSkipsDeletesThatDoNotApply re-inserts id=1 after an equality
delete on it.
+// The delete has a lower sequence number than the new file, so rewriting that
file must keep the row,
+// while the untouched older file must keep its own id=1 hidden.
+func TestCoWRewriteSkipsDeletesThatDoNotApply(t *testing.T) {
+ tbl := newMergeOnReadTestTableVersion(t, "2")
+ tbl = appendRowsOnRef(t, tbl, table.MainBranch, 1, 2, 3)
+ tbl = appendEqualityDelete(t, tbl, []int{1}, `[{"id": 1}]`)
+ tbl = appendRowsOnRef(t, tbl, table.MainBranch, 1, 4, 5)
+ require.Equal(t, []int64{1, 2, 3, 4, 5}, idsInTable(t, tbl))
+
+ tbl = copyOnWriteDelete(t, tbl, 4)
+
+ assert.Equal(t, []int64{1, 2, 3, 5}, idsInTable(t, tbl))
+ assert.Equal(t, int64(5), liveDataRecordCount(t, tbl), "only the file
holding id=4 is rewritten; the other keeps all 3 of its rows")
Review Comment:
this two-file setup is close to what I'd want, but it only rewrites one
file, so it wouldn't catch a mixup inside the new `for _, entry := range
entries` loop in `rewriteScanTasks` (say, applying file A's delete to file B).
Could we add a case that rewrites two files in one operation, each pre-seeded
with a distinct delete on a distinct id, and assert each rewritten file drops
only its own already-deleted row?
While we're here, this scenario only runs through `Delete`; mirroring it
through `Overwrite` the way `TestCoWRewriteKeepsRowsDeleted` does would cover
the filtered-overwrite path through the same code for basically free.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]