zeroshade commented on code in PR #1677:
URL: https://github.com/apache/iceberg-go/pull/1677#discussion_r3798717967


##########
table/snapshot_producers.go:
##########
@@ -953,6 +1001,29 @@ func (sp *snapshotProducer) writeAddedManifest(content 
iceberg.ManifestContent,
        return wr.ToManifestFile(path, counter.Count, 
iceberg.WithManifestFileContent(content))
 }
 
+func (sp *snapshotProducer) writeAddedDeleteManifest(specID int, additions 
[]deleteFileAddition) (_ iceberg.ManifestFile, retErr error) {
+       wr, path, counter, out, err := sp.newManifestWriter(sp.spec(specID), 
iceberg.WithManifestWriterContent(iceberg.ManifestContentDeletes))
+       if err != nil {
+               return nil, err
+       }
+       defer internal.CheckedClose(out, &retErr)
+       defer internal.CheckedClose(wr, &retErr)

Review Comment:
   Could we mirror `writeAddedManifest`'s `writerClosed` guard here? As 
written, the defer closes `wr` again after the explicit `wr.Close()` below, 
relying on `ManifestWriter.Close` remaining idempotent. Matching the adjacent 
lifecycle pattern avoids that assumption and keeps both writer paths consistent.



##########
table/transaction.go:
##########
@@ -1469,42 +1817,202 @@ func (t *Transaction) ReplaceFiles(ctx 
context.Context, dataFilesToDelete, dataF
        // that all files to delete/remove actually exist in the table.
        markedDataForDeletion := make([]iceberg.DataFile, 0, len(setToDelete))
        markedDeleteForRemoval := make([]iceberg.DataFile, 0, 
len(setDeleteFilesToRemove))
+       markedAutoDeleteForRemoval := make([]iceberg.DataFile, 0, 
len(autoSetDeleteFilesToRemove))
        markedDVsForRemoval := make(map[string]iceberg.DataFile, 
len(dvRefsToRemove))
-       for df, err := range s.dataFiles(fs, nil) {
+       removedDeleteSequenceNumbers := make([]int64, 0, 
len(deleteFilesToRemove))
+       removedDeleteContents := make(map[iceberg.ManifestEntryContent]struct{})
+       liveDataFiles := make(map[string]rewriteFileState)
+       survivingPositionDeletes := make([]rewriteFileState, 0)
+       for entry, err := range s.entries(fs, -1) {
                if err != nil {
                        return err
                }
+               df := entry.DataFile()
                path := df.FilePath()
                isData := df.ContentType() == iceberg.EntryContentData
+               isLive := entry.Status() != iceberg.EntryStatusDELETED
+               if isData && isLive {
+                       liveDataFiles[path] = rewriteFileState{file: df, 
dataSequenceNumber: entry.SequenceNum()}
+               }
                if _, ok := setToDelete[path]; ok && isData {
                        markedDataForDeletion = append(markedDataForDeletion, 
df)
                }
                if !isData {
                        if _, ok := setDeleteFilesToRemove[path]; ok {
                                markedDeleteForRemoval = 
append(markedDeleteForRemoval, df)
+                               if seq := entry.SequenceNum(); seq >= 0 {
+                                       removedDeleteSequenceNumbers = 
append(removedDeleteSequenceNumbers, seq)
+                                       removedDeleteContents[df.ContentType()] 
= struct{}{}
+                               } else {
+                                       return fmt.Errorf("delete file %s has 
no data sequence number in the current snapshot", path)
+                               }
+                       } else if _, ok := autoSetDeleteFilesToRemove[path]; ok 
{
+                               markedAutoDeleteForRemoval = 
append(markedAutoDeleteForRemoval, df)
                        } else if ref := 
iceberginternal.BorrowedDataFileReferencedDataFile(df); IsDeletionVector(df) && 
ref != nil {
                                if _, ok := dvRefsToRemove[*ref]; ok {
                                        markedDVsForRemoval[*ref] = df
+                                       if seq := entry.SequenceNum(); seq >= 0 
{
+                                               removedDeleteSequenceNumbers = 
append(removedDeleteSequenceNumbers, seq)
+                                               
removedDeleteContents[df.ContentType()] = struct{}{}
+                                       } else {
+                                               return fmt.Errorf("deletion 
vector %s has no data sequence number in the current snapshot", path)
+                                       }
+                               } else if _, ok := autoDVRefsToRemove[*ref]; ok 
{
+                                       markedDVsForRemoval[*ref] = df
                                }
                        }
                }
                if _, ok := setToAdd[path]; ok {
                        return fmt.Errorf("cannot add files that are already 
referenced by table, files: %s", path)
                }
+               if _, ok := setDeleteFilesToAdd.regularPaths[path]; ok {
+                       return fmt.Errorf("cannot add files that are already 
referenced by table, files: %s", path)
+               }
+               if _, ok := setDeleteFilesToAdd.dvPaths[path]; ok && isLive {
+                       if !IsDeletionVector(df) {
+                               return fmt.Errorf("cannot add deletion vector 
container path already referenced by a non-DV file: %s", path)
+                       }
+                       if offset, length := df.ContentOffset(), 
df.ContentSizeInBytes(); offset != nil && length != nil {
+                               blob := deletionVectorBlobKey{path: path, 
offset: *offset, length: *length}
+                               if _, duplicate := 
setDeleteFilesToAdd.dvBlobs[blob]; duplicate {
+                                       return fmt.Errorf("cannot add deletion 
vector blob already referenced by table: %s at offset %d with length %d",
+                                               path, *offset, *length)
+                               }
+                       }
+               }
+
+               if !isLive || isData {
+                       continue
+               }
+               if IsDeletionVector(df) {
+                       ref := df.ReferencedDataFile()
+                       if ref == nil {
+                               continue
+                       }
+                       if _, addingReplacement := 
setDeleteFilesToAdd.dvsByRef[*ref]; !addingReplacement {
+                               continue
+                       }
+                       if _, explicitlyRemoved := dvRefsToRemove[*ref]; 
explicitlyRemoved {
+                               continue
+                       }
+                       if _, automaticallyRemoved := autoDVRefsToRemove[*ref]; 
automaticallyRemoved {
+                               continue
+                       }
+
+                       return fmt.Errorf("%w: deletion vector for data file %s 
already exists and must be replaced",
+                               ErrInvalidOperation, *ref)
+               }
+               if df.ContentType() == iceberg.EntryContentPosDeletes {
+                       if _, removed := setDeleteFilesToRemove[path]; !removed 
{
+                               if _, automaticallyRemoved := 
autoSetDeleteFilesToRemove[path]; !automaticallyRemoved {
+                                       survivingPositionDeletes = 
append(survivingPositionDeletes, rewriteFileState{
+                                               file: df, dataSequenceNumber: 
entry.SequenceNum(),
+                                       })
+                               }
+                       }
+               }
        }
 
+       addedDataSequenceNumber := nextSequenceNumber
+       if cfg.dataSequenceNumber != nil {
+               addedDataSequenceNumber = *cfg.dataSequenceNumber
+       }
+       for _, df := range dataFilesToAdd {
+               liveDataFiles[df.FilePath()] = rewriteFileState{
+                       file: df, dataSequenceNumber: addedDataSequenceNumber,
+               }
+       }
+       for path := range setToDelete {
+               if _, replacement := setToAdd[path]; !replacement {

Review Comment:
   This replacement branch appears unreachable: the earlier entry scan rejects 
any `setToAdd` path already referenced by the table, so a path cannot 
successfully be present in both `setToDelete` and `setToAdd`. Could this be 
simplified to unconditionally delete every `setToDelete` path from 
`liveDataFiles`?



##########
table/transaction.go:
##########
@@ -1541,7 +2062,14 @@ func (t *Transaction) ReplaceFiles(ctx context.Context, 
dataFilesToDelete, dataF
                return err
        }
 
-       return t.apply(updates, reqs)
+       if err := t.apply(updates, reqs); err != nil {
+               return err
+       }
+       if len(rewriteValidationFiles) > 0 {
+               
t.addValidator(rewriteValidatorWithReferencedDataFiles(rewriteValidationFiles, 
referencedDataFilePaths))

Review Comment:
   For rewrites that add delete files, this extended validator already includes 
a clone of `dataFilesToDelete`, but `RewriteFiles.Commit` subsequently 
registers the basic `rewriteValidator(r.dataFilesToDelete)` as well. That 
validates the same rewritten data files twice and can duplicate manifest 
conflict scanning during retries. Could registration be centralized so exactly 
one validator covers the rewritten files plus any referenced DV targets?



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