laskoviymishka commented on code in PR #2046: URL: https://github.com/apache/iceberg-go/pull/2046#discussion_r4232174267
########## table/cow_rewrite_deletes_test.go: ########## @@ -0,0 +1,391 @@ +// 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 + seededDVs int + }{ + {name: "position delete", formatVersion: "2", seed: deleteID1}, + {name: "deletion vector", formatVersion: "3", seed: deleteID1, seededDVs: 1}, + {name: "equality delete on v2", formatVersion: "2", seed: equalityDeleteID1}, + {name: "equality delete on v3", formatVersion: "3", seed: equalityDeleteID1}, + } + + matchID2 := iceberg.EqualTo(iceberg.Reference("id"), int64(2)) + operations := []cowRewriteOp{ Review Comment: One rewrite case worth adding while we're here: a file whose only surviving row gets removed by the rewrite (say {1,2} with a DV on id=1, then a copy-on-write delete of id=2) exercises the zero-rows-out branch, which none of the current cases hit. Worth asserting the table ends empty, no stray empty data file left live, and `liveDVCount` at 0. ########## table/transaction.go: ########## @@ -2537,7 +2563,7 @@ func (t *Transaction) performMergeOnReadDeletion(ctx context.Context, snapshotPr updater := t.updateSnapshot(wfs, snapshotProps, OpDelete).mergeOverwrite(&commitUUID, filter) updater.setManifestConcurrency(concurrency) - filesToDelete, withPartialDeletions, _, err := t.classifyFilesForDeletions(ctx, fs, filter, caseSensitive, concurrency) + filesToDelete, withPartialDeletions, err := t.classifyFilesForDeletions(ctx, fs, filter, caseSensitive, concurrency) Review Comment: The copy-on-write drop and rewrite paths both route through `planningDeletes.removeDataFile` now, which also drops the DV. The merge-on-read path right below (the unchanged `for _, df := range filesToDelete { updater.deleteDataFile(df) }` loop) still bare-deletes, so on a v3 table a MoR delete that fully matches a file carrying a DV leaves that vector live and orphaned. That's the exact gap we closed on copy-on-write last round, just on the sibling path. I'd apply the same gate here: when `meta.formatVersion >= 3 && len(filesToDelete) > 0`, read the planning deletes and route the drop loop through `removeDataFile`. If you'd rather scope it out of this PR, that's fine, but let's say so in the description and file the follow-up, and add a MoR whole-file-drop case alongside `TestCoWDeleteRemovesDeletionVectorsOfDroppedFiles` to pin it. ########## table/transaction.go: ########## @@ -2833,17 +2845,75 @@ 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 +} + +// planningDeletes indexes the live delete files of the planning snapshot the same way scan planning indexes them. +type planningDeletes struct { + positional *positionalDeleteIndex + dvs map[string]iceberg.ManifestEntry + equality *equalityDeleteIndex +} + +// readPlanningDeletes reads and indexes every live delete file of the planning snapshot. +// +// 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 (#2160): 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) readPlanningDeletes(fs io.IO, builtMeta Metadata) (planningDeletes, error) { + meta, err := t.txnMeta() + if err != nil { + return planningDeletes{}, err + } + + var liveDeletes []iceberg.ManifestEntry + if s := t.planningSnapshot(meta); s != nil { + for entry, err := range s.entries(fs, iceberg.ManifestContentDeletes, true) { + if err != nil { + return planningDeletes{}, fmt.Errorf("failed to read delete manifests: %w", err) + } + liveDeletes = append(liveDeletes, entry) + } + } + + classified, err := classifyManifestEntries(liveDeletes) + if err != nil { + return planningDeletes{}, err + } + positional, err := buildPositionalDeleteIndex(classified.positionalDeleteEntries) + if err != nil { + return planningDeletes{}, err + } + dvs, err := buildDVIndex(classified.dvEntries) + if err != nil { + return planningDeletes{}, err + } + equality, err := buildEqualityDeleteIndex(classified.equalityDeleteEntries, builtMeta, meta.CurrentSchema()) + if err != nil { + return planningDeletes{}, err + } + + return planningDeletes{positional: positional, dvs: dvs, equality: equality}, nil +} + +// removeDataFile removes df together with the deletion vector that references it. +// Position and equality delete files are left alone, since they can still apply to other files. +func (d planningDeletes) removeDataFile(updater *snapshotProducer, df iceberg.DataFile) { + updater.deleteDataFile(df) + if dv, ok := d.dvs[df.FilePath()]; ok { + updater.removeDeletionVector(dv.DataFile()) Review Comment: Routing the DV removal through here flips the whole commit to `noReplay` (`removeDeletionVector` sets `noReplay = true` in snapshot_producers.go), so a copy-on-write delete/overwrite that drops or rewrites a DV-bearing file no longer refreshes-and-replays on a CAS conflict. An unrelated concurrent append now surfaces `ErrCommitFailed` where it used to retry. Stricter-than-Java is a defensible call, but it's a behavior change for streaming/CDC callers that isn't noted anywhere. I'd call it out in the `Delete`/`Overwrite` docs and the PR description, and a test expecting `ErrCommitFailed` on concurrent-append-then-copy-on-write-delete would lock the intent in. -- 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]
