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]

Reply via email to