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]