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]