laskoviymishka commented on code in PR #2046:
URL: https://github.com/apache/iceberg-go/pull/2046#discussion_r4145292294


##########
table/transaction.go:
##########
@@ -2807,17 +2792,73 @@ 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.
+//
+// Known limitation: every live delete manifest of the snapshot is read on 
each call,
+// with no partition pruning, so a narrow rewrite on a delete-heavy table pays 
for all of them.
+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

Review Comment:
   This bare `return nil, err` is the same `meta.Build()` failure that 
`rewriteSingleFile` wraps as `failed to build metadata: %w` a couple calls 
down; worth matching the wrap so triage is consistent across the one pipeline. 
Not blocking.



##########
table/transaction.go:
##########
@@ -2849,10 +2894,10 @@ func (t *Transaction) rewriteFilesWithFilter(ctx 
context.Context, fs io.IO, upda
                }
                rewrittenFiles, err := t.rewriteSingleFile(ctx, args)
                if err != nil {
-                       return fmt.Errorf("failed to rewrite file %s: %w", 
originalFile.FilePath(), err)
+                       return fmt.Errorf("failed to rewrite file %s: %w", 
task.File.FilePath(), err)
                }
 
-               updater.deleteDataFile(originalFile)
+               updater.deleteDataFile(task.File)

Review Comment:
   When we drop `task.File` here we also need to retire the deletion vector 
that pointed at it. The spec is explicit: *"When removing a data file, writers 
must also remove any deletion vector that applies to that data file from delete 
manifests."* Right now nothing does, so the DV stays live in the delete 
manifest.
   
   It doesn't corrupt reads (the rewritten file gets a fresh path, so the 
orphaned DV never matches anything again), but the Puffin blob is now 
unreachable by expire/orphan cleanup, and the snapshot summary's 
`removed-dvs`/`total-delete-files` diverge from what Java writes for the 
identical commit. Java does this inline on every commit 
(`MergingSnapshotProducer.apply` → `removeDanglingDeletesFor`), which 
`BaseOverwriteFiles` (the class this path mirrors) inherits, so it's not a 
deferred maintenance step there.
   
   `fileScanTaskForDataEntry` already matched and populated the DVs on the 
task, so this is close to free:
   
   ```go
   updater.deleteDataFile(task.File)
   for _, dv := range task.DeletionVectorFiles {
       updater.removeDeletionVector(dv)
   }
   ```
   
   The full-match `filesToDelete` loop up in `performCopyOnWriteDeletion` has 
the same gap (pre-existing there, not introduced here), so the fix should cover 
both loops. Keep it DV-scoped: Java's `isDanglingDV` deliberately leaves v2 
position/equality delete files alone, and the equality/position subtests 
already match that. Worth pinning with a `liveDVCount` assertion after the 
rewrite in the deletion-vector cases, which today only check it before.



##########
table/transaction.go:
##########
@@ -2807,17 +2792,73 @@ 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.
+//
+// Known limitation: every live delete manifest of the snapshot is read on 
each call,
+// with no partition pruning, so a narrow rewrite on a delete-heavy table pays 
for all of them.
+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()

Review Comment:
   I'd build `builtMeta` once here and thread it onto `rewriteSingleFileArgs` 
for reuse, rather than letting `rewriteSingleFile` call `meta.Build()` again 
per file. Beyond the repeated rebuild, `Build()` isn't idempotent on a live 
builder: when `HasChanges()` is true, `buildCommonMetadata` appends 
`previousFileEntry` to `metadataLog` on every call and never clears it. So if 
an earlier step this transaction already staged an update (the auto 
name-mapping `SetProperties` before classification, say), each extra `Build()` 
writes a duplicate previous-files entry into persisted table history, not just 
a transient allocation.



##########
table/cow_rewrite_deletes_test.go:
##########
@@ -0,0 +1,288 @@
+// 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/arrow-go/v18/arrow/array"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/iceberg-go"
+       "github.com/apache/iceberg-go/table"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+)
+
+type cowRewriteOp struct {
+       name        string
+       run         func(*testing.T, *table.Table) *table.Table
+       wantIDs     []int64
+       wantRecords int64
+}
+
+// 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()
+
+               return mergeOnReadDelete(t, tbl, 1)
+       }
+       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},

Review Comment:
   Each case seeds exactly one delete kind onto the rewritten file, but 
production builds all three 
(`DeleteFiles`/`EqualityDeleteFiles`/`DeletionVectorFiles`) on the same task 
independently. A case that stacks an equality delete and a DV on different rows 
of one file, asserting both drop, would pin the combined path against a future 
change that early-returns after matching just one kind. Not blocking.



-- 
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