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


##########
table/equality_delete_reader.go:
##########
@@ -44,6 +45,9 @@ import (
 
 var ErrAmbiguousEqualityColumn = errors.New("equality delete column is 
ambiguous")
 
+// ErrConflictingEqualityDeleteMetadata indicates incompatible metadata for 
one delete-file path.
+var ErrConflictingEqualityDeleteMetadata = errors.New("conflicting equality 
delete metadata")

Review Comment:
   Small carryover: I'd have this doc comment name the two things that actually 
trigger it (a different equality field ID set, or a different file format for 
the same path) so a downstream consumer can classify what it is looking at. 
Something like "returned when the same delete-file path appears with 
incompatible metadata: either a different equality field ID set or a different 
file format."



##########
table/equality_delete_reader.go:
##########
@@ -364,40 +423,101 @@ func newLazyEqualityDeleteLoader(
                tableSchema:  tableSchema,
                tableSchemas: tableSchemas,
                nameMapping:  nameMapping,
-               files:        make(map[string]*lazyEqualityDeleteFile),
        }
 
+       var firstPath string
+       var firstFile *lazyEqualityDeleteFile
        for _, task := range tasks {
                for _, dataFile := range task.EqualityDeleteFiles {
+                       if loader.files == nil && firstFile != nil &&
+                               firstFile.hasPointerIdentity && dataFile == 
firstFile.dataFile {

Review Comment:
   This pointer-identity skip runs before the `ContentType()` check just below 
it, so it short-circuits without confirming the file is an eq-delete entry. It 
is safe today, since `firstFile.dataFile` already cleared ContentType and 
pointer identity means it is the same object, but the invariant is implicit and 
it flips the check-before-skip order used everywhere else in this loop. I'd 
move it just below the ContentType guard so the ordering matches.



##########
table/equality_delete_reader.go:
##########
@@ -303,6 +307,59 @@ func newEqualityDeleteFileSet(id int, deleteSet 
*equalityDeleteSet) *equalityDel
        }
 }
 
+func sameEqualityFieldIDSet(left, right []int) bool {

Review Comment:
   I'd swap the two `slices.Contains` passes for sorting copies of each side 
and comparing with `slices.Equal`. As written this treats the IDs as a set, so 
`[1,1,2]` and `[1,2,2]` compare equal even though the multisets differ, and the 
spec types equality_ids as `list<int>` (not `set`), so duplicate IDs can show 
up in corrupt metadata and we'd silently drop the conflict. Sorting handles the 
multiset case correctly and lets the second loop go, since it is only ever 
reachable for duplicate IDs anyway.
   
   ```go
   func sameEqualityFieldIDSet(left, right []int) bool {
        if len(left) != len(right) {
                return false
        }
        l, r := slices.Clone(left), slices.Clone(right)
        slices.Sort(l)
        slices.Sort(r)
        return slices.Equal(l, r)
   }
   ```
   
   While we're here, the "set predicate" comment at the call site reads a touch 
stronger than the spec. It is not that the IDs are a spec set-type; it is that 
equality matching is order-independent (an AND over the column equalities), so 
the same IDs in a different order name the same predicate. Worth a word tweak 
so it doesn't read as a spec claim.



##########
table/equality_delete_metadata_validation_test.go:
##########
@@ -0,0 +1,176 @@
+// 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
+
+import (
+       "bytes"
+       "fmt"
+       "slices"
+       "testing"
+
+       "github.com/apache/iceberg-go"
+       iceio "github.com/apache/iceberg-go/io"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+)
+
+type nonComparableEqualityMetadataFile struct {
+       iceberg.DataFile
+       fieldIDs []int
+}
+
+func (f nonComparableEqualityMetadataFile) EqualityFieldIDs() []int {
+       return slices.Clone(f.fieldIDs)
+}
+
+func TestEqualityDeleteMetadataAcceptsDistinctSamePathFiles(t *testing.T) {
+       t.Parallel()
+
+       schema := iceberg.NewSchema(0,
+               iceberg.NestedField{ID: 1, Name: "id", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+       )
+       wrappers := []struct {
+               name string
+               wrap func(iceberg.DataFile) iceberg.DataFile
+       }{
+               {
+                       name: "borrowed metadata",
+                       wrap: func(file iceberg.DataFile) iceberg.DataFile { 
return file },
+               },
+               {
+                       name: "public getter with non-comparable interface 
value",
+                       wrap: func(file iceberg.DataFile) iceberg.DataFile {
+                               // The outer struct is comparable, but its 
interface value is not.
+                               return struct{ iceberg.DataFile 
}{nonComparableEqualityMetadataFile{
+                                       DataFile: file,
+                                       fieldIDs: file.EqualityFieldIDs(),
+                               }}
+                       },
+               },
+       }
+
+       for _, wrapper := range wrappers {
+               for _, useMap := range []bool{false, true} {
+                       t.Run(fmt.Sprintf("%s/map=%t", wrapper.name, useMap), 
func(t *testing.T) {
+                               fs := &countingEqualityDeleteOpenFS{MemFS: 
iceio.NewMemFS()}
+                               path := "mem://metadata-accept/delete.parquet"
+                               writeEqualityDeleteParquetToMemFS(t, fs.MemFS, 
path, `[{"id": 1}]`)
+                               first := 
newEqualityDeleteSetAssemblyTestFile(t, path, []int{1})
+                               second := 
newEqualityDeleteSetAssemblyTestFile(t, path, []int{1})
+                               require.NotSame(t, first, second)
+                               tasks := []FileScanTask{{EqualityDeleteFiles: 
[]iceberg.DataFile{wrapper.wrap(first)}}}
+                               uniqueFiles := 1
+                               if useMap {
+                                       otherPath := 
"mem://metadata-accept/other.parquet"
+                                       writeEqualityDeleteParquetToMemFS(t, 
fs.MemFS, otherPath, `[{"id": 2}]`)
+                                       other := 
newEqualityDeleteSetAssemblyTestFile(t, otherPath, []int{1})
+                                       tasks = append(tasks, 
FileScanTask{EqualityDeleteFiles: []iceberg.DataFile{other}})
+                                       uniqueFiles++
+                               }
+                               tasks = append(tasks, 
FileScanTask{EqualityDeleteFiles: []iceberg.DataFile{wrapper.wrap(second)}})
+
+                               loader, err := newLazyEqualityDeleteLoader(fs, 
schema, nil, nil, tasks)
+                               require.NoError(t, err)
+                               require.NotNil(t, loader)
+                               assert.Len(t, loader.files, uniqueFiles)
+                               assert.Zero(t, fs.attempts.Load())
+                               var firstSet *equalityDeleteSet
+                               for i, task := range tasks {
+                                       sets, err := loader.load(t.Context(), 
task)
+                                       require.NoError(t, err)
+                                       require.Len(t, sets, 1)
+                                       if i == 0 {
+                                               firstSet = sets[0]
+                                       } else if i == len(tasks)-1 {
+                                               assert.Same(t, firstSet, 
sets[0])
+                                       }
+                               }
+                               assert.Equal(t, int64(uniqueFiles), 
fs.attempts.Load())
+                               assert.Equal(t, int64(uniqueFiles), 
fs.opens.Load())
+                               var key bytes.Buffer
+                               key.WriteByte(1)
+                               bufPutUint64(&key, 1)
+                               assert.Equal(t, set[string]{key.String(): {}}, 
firstSet.keys)
+                               assert.Equal(t, []int{1}, firstSet.fieldIDs)
+
+                               fs.attempts.Store(0)
+                               fs.opens.Store(0)
+                               perTask, err := 
readAllEqualityDeleteFiles(t.Context(), fs, schema, nil, tasks, 1)
+                               require.NoError(t, err)
+                               require.Len(t, perTask, len(tasks))
+                               require.Len(t, perTask[0], 1)
+                               require.Len(t, perTask[len(tasks)-1], 1)
+                               assert.Same(t, perTask[0][0], 
perTask[len(tasks)-1][0])
+                               assert.Equal(t, firstSet, perTask[0][0])
+                               assert.Equal(t, int64(uniqueFiles), 
fs.attempts.Load())
+                               assert.Equal(t, int64(uniqueFiles), 
fs.opens.Load())
+                       })
+               }
+       }
+}
+
+func TestEqualityDeleteMetadataConflictErrors(t *testing.T) {
+       t.Parallel()
+
+       schema := iceberg.NewSchema(0,
+               iceberg.NestedField{ID: 1, Name: "id", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+               iceberg.NestedField{ID: 2, Name: "data", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+       )
+       path := "mem://metadata-conflict/delete.parquet"
+       first := newEqualityDeleteSetAssemblyTestFile(t, path, []int{1, 2})
+       for _, tt := range []struct {
+               name     string
+               fieldIDs []int
+               format   iceberg.FileFormat
+               wantErr  error
+       }{
+               {"different IDs", []int{1, 3}, iceberg.ParquetFile, 
ErrConflictingEqualityDeleteMetadata},
+               {"reordered IDs", []int{2, 1}, iceberg.ParquetFile, nil},
+               {"different format", []int{1, 2}, iceberg.AvroFile, 
ErrConflictingEqualityDeleteMetadata},
+               {"empty IDs", nil, iceberg.ParquetFile, 
ErrEmptyEqualityFieldIDs},
+       } {
+               t.Run(tt.name, func(t *testing.T) {
+                       builder, err := iceberg.NewDataFileBuilder(
+                               *iceberg.UnpartitionedSpec, 
iceberg.EntryContentEqDeletes, path,
+                               tt.format, nil, nil, nil, 1, 128)
+                       require.NoError(t, err)
+                       second := builder.EqualityFieldIDs(tt.fieldIDs).Build()
+                       tasks := []FileScanTask{
+                               {EqualityDeleteFiles: 
[]iceberg.DataFile{first}},
+                               {EqualityDeleteFiles: 
[]iceberg.DataFile{second}},
+                       }
+                       fs := &countingEqualityDeleteOpenFS{MemFS: 
iceio.NewMemFS()}
+                       loader, err := newLazyEqualityDeleteLoader(fs, schema, 
nil, nil, tasks)
+                       if tt.wantErr == nil {

Review Comment:
   The `wantErr == nil` branch returns here before the 
`readAllEqualityDeleteFiles` call lower down, so the eager path never sees the 
reordered-IDs accept case; it is only checked against 
`newLazyEqualityDeleteLoader`. Since flipping reorder to accept is the behavior 
this test is pinning, I'd add a `readAllEqualityDeleteFiles` NoError assertion 
in this branch before the return so both paths are covered.



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