This is an automated email from the ASF dual-hosted git repository.

laskoviymishka pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-go.git


The following commit(s) were added to refs/heads/main by this push:
     new b8ee7dd64 perf(table): deduplicate equality delete metadata setup 
(#2025)
b8ee7dd64 is described below

commit b8ee7dd647dbe6b5e66dbdfb58f40519399a1104
Author: Minh Vu <[email protected]>
AuthorDate: Mon Oct 5 21:13:06 2026 +0200

    perf(table): deduplicate equality delete metadata setup (#2025)
---
 table/arrow_scanner.go                            |  19 +-
 table/data_file_stats_ref.go                      |  27 ++
 table/data_file_stats_ref_test.go                 |  11 +
 table/equality_delete_metadata_validation_test.go | 176 +++++++++++++
 table/equality_delete_reader.go                   | 286 ++++++++++++++++------
 table/equality_delete_reader_bench_test.go        |  78 +++++-
 table/equality_delete_reader_internal_test.go     |  93 +++++++
 7 files changed, 599 insertions(+), 91 deletions(-)

diff --git a/table/arrow_scanner.go b/table/arrow_scanner.go
index 4d29a69e4..c7fe27198 100644
--- a/table/arrow_scanner.go
+++ b/table/arrow_scanner.go
@@ -2629,25 +2629,14 @@ func (as *arrowScan) GetRecords(ctx context.Context, 
tasks []FileScanTask) (*arr
                return nil, nil, err
        }
 
-       var tableSchemas []*iceberg.Schema
-
-loadSchemaHistory:
-       for _, task := range tasks {
-               for _, deleteFile := range task.EqualityDeleteFiles {
-                       for _, fieldID := range deleteFile.EqualityFieldIDs() {
-                               if _, found := 
invariants.tableSchema.FindFieldByID(fieldID); !found {
-                                       tableSchemas = as.metadata.Schemas()
-
-                                       break loadSchemaHistory
-                               }
-                       }
-               }
-       }
        equalityDeleteLoader, err := newLazyEqualityDeleteLoader(
-               as.fs, invariants.tableSchema, tableSchemas, 
invariants.nameMapping, tasks)
+               as.fs, invariants.tableSchema, nil, invariants.nameMapping, 
tasks)
        if err != nil {
                return nil, nil, err
        }
+       if equalityDeleteLoader.needsSchemaHistory() {
+               equalityDeleteLoader.tableSchemas = as.metadata.Schemas()
+       }
        equalityDeleteLoader.addFieldIDs(invariants.projectedIDs)
 
        positionDeleteLoader := newLazyPositionDeleteLoader(as.fs, tasks)
diff --git a/table/data_file_stats_ref.go b/table/data_file_stats_ref.go
index 8d4c2fa2c..bce006107 100644
--- a/table/data_file_stats_ref.go
+++ b/table/data_file_stats_ref.go
@@ -18,6 +18,8 @@
 package table
 
 import (
+       "slices"
+
        iceberg "github.com/apache/iceberg-go"
        "github.com/apache/iceberg-go/internal"
 )
@@ -45,6 +47,31 @@ func dataFileCollections(file iceberg.DataFile) (
        return internal.BorrowedDataFileCollections(file)
 }
 
+// dataFileEqualityFieldIDsRef returns a borrowed equality field ID view for
+// short-lived comparisons. Callers must not retain or mutate the returned 
slice.
+func dataFileEqualityFieldIDsRef(file iceberg.DataFile) []int {
+       if ref, ok := file.(internal.DataFileCollectionsRef); ok {
+               _, _, _, equalityFieldIDs := 
ref.DataFileCollectionsRef(internal.DataFileRef{})
+
+               return equalityFieldIDs
+       }
+
+       return file.EqualityFieldIDs()
+}
+
+// dataFileEqualityFieldIDs returns equality field IDs that are safe to retain.
+// The internal collection view is borrowed, so clone it before storing the IDs
+// on scan-lifetime loader state. The public getter already returns its own 
view.
+func dataFileEqualityFieldIDs(file iceberg.DataFile) []int {
+       if ref, ok := file.(internal.DataFileCollectionsRef); ok {
+               _, _, _, equalityFieldIDs := 
ref.DataFileCollectionsRef(internal.DataFileRef{})
+
+               return slices.Clone(equalityFieldIDs)
+       }
+
+       return file.EqualityFieldIDs()
+}
+
 // dataFilePartition returns a borrowed partition map for the concrete
 // manifest data file and falls back to the public getter for other DataFile
 // implementations. Callers must use the map only for the current planning
diff --git a/table/data_file_stats_ref_test.go 
b/table/data_file_stats_ref_test.go
index 9b83d5129..e5e1dee42 100644
--- a/table/data_file_stats_ref_test.go
+++ b/table/data_file_stats_ref_test.go
@@ -183,6 +183,17 @@ func TestDataFileCollectionsUsesBorrowedView(t *testing.T) 
{
        assert.Equal(t, map[int]int64{1: 10}, measuredColumnSizes)
 }
 
+func TestDataFileEqualityFieldIDsClonesBorrowedView(t *testing.T) {
+       file := testDataFileWithStats(t)
+       require.Implements(t, (*internal.DataFileCollectionsRef)(nil), file)
+
+       fieldIDs := dataFileEqualityFieldIDs(file)
+       require.Equal(t, []int{1}, fieldIDs)
+       fieldIDs[0] = 99
+
+       assert.Equal(t, []int{1}, file.EqualityFieldIDs())
+}
+
 func TestDataFileCollectionsFallsBackToPublicGetters(t *testing.T) {
        file := &publicStatsDataFile{DataFile: testDataFileWithStats(t)}
 
diff --git a/table/equality_delete_metadata_validation_test.go 
b/table/equality_delete_metadata_validation_test.go
new file mode 100644
index 000000000..a180fc277
--- /dev/null
+++ b/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 {
+                               require.NoError(t, err)
+                               require.NotNil(t, loader)
+                               assert.Len(t, loader.files, 1)
+                               assert.Zero(t, fs.attempts.Load())
+
+                               return
+                       }
+
+                       require.ErrorIs(t, err, tt.wantErr)
+                       require.ErrorContains(t, err, path)
+                       _, err = readAllEqualityDeleteFiles(t.Context(), fs, 
schema, nil, tasks, 1)
+                       require.ErrorIs(t, err, tt.wantErr)
+                       require.ErrorContains(t, err, path)
+                       assert.Zero(t, fs.attempts.Load(), "reject conflicting 
metadata before I/O")
+               })
+       }
+}
diff --git a/table/equality_delete_reader.go b/table/equality_delete_reader.go
index 17993313b..400965701 100644
--- a/table/equality_delete_reader.go
+++ b/table/equality_delete_reader.go
@@ -26,6 +26,7 @@ import (
        "fmt"
        "maps"
        "math"
+       "reflect"
        "slices"
        "sync"
        "unsafe"
@@ -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")
+
 // equalityDeleteSet holds the set of delete keys and the column names
 // used to look them up in data records. Each set corresponds to one
 // group of equality field IDs — delete files with different field IDs
@@ -303,6 +307,59 @@ func newEqualityDeleteFileSet(id int, deleteSet 
*equalityDeleteSet) *equalityDel
        }
 }
 
+func sameEqualityFieldIDSet(left, right []int) bool {
+       if len(left) != len(right) {
+               return false
+       }
+
+       for _, id := range left {
+               if !slices.Contains(right, id) {
+                       return false
+               }
+       }
+       for _, id := range right {
+               if !slices.Contains(left, id) {
+                       return false
+               }
+       }
+
+       return true
+}
+
+func dataFileHasPointerIdentity(dataFile iceberg.DataFile) bool {
+       return reflect.TypeOf(dataFile).Kind() == reflect.Pointer
+}
+
+func validateEqualityDeleteMetadata(
+       dataFile iceberg.DataFile,
+       existingFile iceberg.DataFile,
+       existingFieldIDs []int,
+       existingHasPointerIdentity bool,
+) error {
+       // Callers have already inspected ContentType, so both files must be 
non-nil.
+       // A pointer-backed retained file makes interface identity comparison 
safe,
+       // even when a later custom DataFile value itself is not comparable.
+       if existingHasPointerIdentity && dataFile == existingFile {
+               return nil
+       }
+
+       fieldIDs := dataFileEqualityFieldIDsRef(dataFile)
+       if len(fieldIDs) == 0 {
+               return fmt.Errorf("%w: equality delete file %s", 
ErrEmptyEqualityFieldIDs, dataFile.FilePath())
+       }
+       // Equality delete field IDs are a set predicate. The first file keeps 
its
+       // encoding order; later entries for the same path may list the same 
IDs in
+       // a different order.
+       if dataFile.FileFormat() == existingFile.FileFormat() && 
sameEqualityFieldIDSet(fieldIDs, existingFieldIDs) {
+               return nil
+       }
+
+       return fmt.Errorf(
+               "%w for file %s: first format=%s equality field IDs=%v, later 
format=%s equality field IDs=%v",
+               ErrConflictingEqualityDeleteMetadata, dataFile.FilePath(),
+               existingFile.FileFormat(), existingFieldIDs, 
dataFile.FileFormat(), fieldIDs)
+}
+
 func schemaForEqualityFields(current *iceberg.Schema, schemas 
[]*iceberg.Schema, fieldIDs []int) *iceberg.Schema {
        hasAllFields := func(schema *iceberg.Schema) bool {
                for _, fieldID := range fieldIDs {
@@ -334,13 +391,15 @@ type lazyEqualityDeleteLoader struct {
        tableSchemas []*iceberg.Schema
        nameMapping  iceberg.NameMapping
        files        map[string]*lazyEqualityDeleteFile
+       singleFile   *lazyEqualityDeleteFile
        combinations sync.Map
 }
 
 type lazyEqualityDeleteFile struct {
-       id       int
-       dataFile iceberg.DataFile
-       fieldIDs []int
+       id                 int
+       dataFile           iceberg.DataFile
+       fieldIDs           []int
+       hasPointerIdentity bool
 
        once sync.Once
        set  *equalityDeleteFileSet
@@ -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 {
+                               continue
+                       }
                        if dataFile.ContentType() != 
iceberg.EntryContentEqDeletes {
                                continue
                        }
 
-                       fieldIDs := dataFile.EqualityFieldIDs()
-                       if len(fieldIDs) == 0 {
-                               return nil, fmt.Errorf("%w: equality delete 
file %s", ErrEmptyEqualityFieldIDs, dataFile.FilePath())
+                       path := dataFile.FilePath()
+                       if loader.files == nil {
+                               if firstFile == nil {
+                                       fieldIDs := 
dataFileEqualityFieldIDs(dataFile)
+                                       if len(fieldIDs) == 0 {
+                                               return nil, fmt.Errorf("%w: 
equality delete file %s", ErrEmptyEqualityFieldIDs, path)
+                                       }
+
+                                       firstPath = path
+                                       firstFile = &lazyEqualityDeleteFile{
+                                               dataFile:           dataFile,
+                                               fieldIDs:           fieldIDs,
+                                               hasPointerIdentity: 
dataFileHasPointerIdentity(dataFile),
+                                       }
+
+                                       continue
+                               }
+                               if path == firstPath {
+                                       if err := 
validateEqualityDeleteMetadata(
+                                               dataFile, firstFile.dataFile, 
firstFile.fieldIDs, firstFile.hasPointerIdentity,
+                                       ); err != nil {
+                                               return nil, err
+                                       }
+
+                                       continue
+                               }
+
+                               loader.files = 
make(map[string]*lazyEqualityDeleteFile, 2)
+                               loader.files[firstPath] = firstFile
                        }
+                       if file, ok := loader.files[path]; ok {
+                               if err := validateEqualityDeleteMetadata(
+                                       dataFile, file.dataFile, file.fieldIDs, 
file.hasPointerIdentity,
+                               ); err != nil {
+                                       return nil, err
+                               }
 
-                       path := dataFile.FilePath()
-                       if _, ok := loader.files[path]; ok {
                                continue
                        }
 
+                       fieldIDs := dataFileEqualityFieldIDs(dataFile)
+                       if len(fieldIDs) == 0 {
+                               return nil, fmt.Errorf("%w: equality delete 
file %s", ErrEmptyEqualityFieldIDs, path)
+                       }
+
                        loader.files[path] = &lazyEqualityDeleteFile{
-                               id:       len(loader.files),
-                               dataFile: dataFile,
-                               fieldIDs: fieldIDs,
+                               id:                 len(loader.files),
+                               dataFile:           dataFile,
+                               fieldIDs:           fieldIDs,
+                               hasPointerIdentity: 
dataFileHasPointerIdentity(dataFile),
                        }
                }
        }
 
-       if len(loader.files) == 0 {
+       if firstFile == nil {
                return nil, nil
        }
+       if loader.files == nil {
+               loader.files = map[string]*lazyEqualityDeleteFile{firstPath: 
firstFile}
+               loader.singleFile = firstFile
+       }
 
        return loader, nil
 }
 
+func (l *lazyEqualityDeleteLoader) needsSchemaHistory() bool {
+       if l == nil {
+               return false
+       }
+
+       for _, file := range l.files {
+               for _, fieldID := range file.fieldIDs {
+                       if _, found := l.tableSchema.FindFieldByID(fieldID); 
!found {
+                               return true
+                       }
+               }
+       }
+
+       return false
+}
+
 func (l *lazyEqualityDeleteLoader) addFieldIDs(idset set[int]) {
        if l == nil {
                return
@@ -412,7 +532,10 @@ func (l *lazyEqualityDeleteLoader) addFieldIDs(idset 
set[int]) {
 
 func (l *lazyEqualityDeleteLoader) loadFile(ctx context.Context, file 
*lazyEqualityDeleteFile) (*equalityDeleteFileSet, error) {
        file.once.Do(func() {
-               deleteSchema := schemaForEqualityFields(l.tableSchema, 
l.tableSchemas, file.fieldIDs)
+               deleteSchema := l.tableSchema
+               if len(l.tableSchemas) > 0 {
+                       deleteSchema = schemaForEqualityFields(l.tableSchema, 
l.tableSchemas, file.fieldIDs)
+               }
                keys, colNames, err := readEqualityDeleteFile(
                        ctx, l.fs, deleteSchema, l.nameMapping, file.dataFile, 
file.fieldIDs)
                if err != nil {
@@ -453,13 +576,17 @@ func (l *lazyEqualityDeleteLoader) load(ctx 
context.Context, task FileScanTask)
        }
        if len(task.EqualityDeleteFiles) == 1 {
                dataFile := task.EqualityDeleteFiles[0]
-               if dataFile.ContentType() != iceberg.EntryContentEqDeletes {
-                       return nil, nil
-               }
+               file := l.singleFile
+               if file == nil || !file.hasPointerIdentity || dataFile != 
file.dataFile {
+                       if dataFile.ContentType() != 
iceberg.EntryContentEqDeletes {
+                               return nil, nil
+                       }
 
-               file, ok := l.files[dataFile.FilePath()]
-               if !ok {
-                       return nil, nil
+                       var ok bool
+                       file, ok = l.files[dataFile.FilePath()]
+                       if !ok {
+                               return nil, nil
+                       }
                }
 
                fileSet, err := l.loadFile(ctx, file)
@@ -473,18 +600,13 @@ func (l *lazyEqualityDeleteLoader) load(ctx 
context.Context, task FileScanTask)
                return []*equalityDeleteSet{fileSet.equalityDeleteSet}, nil
        }
 
-       perFile := make(map[string]*equalityDeleteFileSet, 
len(task.EqualityDeleteFiles))
+       files := make([]*equalityDeleteFileSet, 0, 
len(task.EqualityDeleteFiles))
        for _, dataFile := range task.EqualityDeleteFiles {
                if dataFile.ContentType() != iceberg.EntryContentEqDeletes {
                        continue
                }
 
-               path := dataFile.FilePath()
-               if _, seen := perFile[path]; seen {
-                       continue
-               }
-
-               file, ok := l.files[path]
+               file, ok := l.files[dataFile.FilePath()]
                if !ok {
                        continue
                }
@@ -493,14 +615,10 @@ func (l *lazyEqualityDeleteLoader) load(ctx 
context.Context, task FileScanTask)
                if err != nil {
                        return nil, err
                }
-               perFile[path] = fileSet
+               files = append(files, fileSet)
        }
 
-       if len(perFile) == 0 {
-               return nil, nil
-       }
-
-       return buildEqualityDeleteSetsForTask(task, perFile, l.combine), nil
+       return buildEqualityDeleteSetsForFiles(files, l.combine), nil
 }
 
 // readAllEqualityDeleteFiles reads all unique equality delete files from
@@ -509,13 +627,13 @@ func (l *lazyEqualityDeleteLoader) load(ctx 
context.Context, task FileScanTask)
 // kept as separate sets (not merged).
 func readAllEqualityDeleteFiles(ctx context.Context, fs iceio.IO, schema 
*iceberg.Schema, nameMapping iceberg.NameMapping, tasks []FileScanTask, 
concurrency int) (map[int][]*equalityDeleteSet, error) {
        type deleteFileInfo struct {
-               id       int
-               file     iceberg.DataFile
-               fieldIDs []int
+               id                 int
+               file               iceberg.DataFile
+               fieldIDs           []int
+               hasPointerIdentity bool
        }
 
        uniqueDeletes := make(map[string]deleteFileInfo)
-       hasAny := false
 
        for _, t := range tasks {
                for _, d := range t.EqualityDeleteFiles {
@@ -523,22 +641,32 @@ func readAllEqualityDeleteFiles(ctx context.Context, fs 
iceio.IO, schema *iceber
                                continue
                        }
 
-                       if len(d.EqualityFieldIDs()) == 0 {
-                               return nil, fmt.Errorf("%w: equality delete 
file %s", ErrEmptyEqualityFieldIDs, d.FilePath())
+                       path := d.FilePath()
+                       if info, ok := uniqueDeletes[path]; ok {
+                               if err := validateEqualityDeleteMetadata(
+                                       d, info.file, info.fieldIDs, 
info.hasPointerIdentity,
+                               ); err != nil {
+                                       return nil, err
+                               }
+
+                               continue
                        }
 
-                       hasAny = true
-                       if _, ok := uniqueDeletes[d.FilePath()]; !ok {
-                               uniqueDeletes[d.FilePath()] = deleteFileInfo{
-                                       id:       len(uniqueDeletes),
-                                       file:     d,
-                                       fieldIDs: d.EqualityFieldIDs(),
-                               }
+                       fieldIDs := dataFileEqualityFieldIDs(d)
+                       if len(fieldIDs) == 0 {
+                               return nil, fmt.Errorf("%w: equality delete 
file %s", ErrEmptyEqualityFieldIDs, path)
+                       }
+
+                       uniqueDeletes[path] = deleteFileInfo{
+                               id:                 len(uniqueDeletes),
+                               file:               d,
+                               fieldIDs:           fieldIDs,
+                               hasPointerIdentity: 
dataFileHasPointerIdentity(d),
                        }
                }
        }
 
-       if !hasAny {
+       if len(uniqueDeletes) == 0 {
                return nil, nil
        }
 
@@ -637,39 +765,42 @@ func buildEqualityDeleteSetsForTask(
                return []*equalityDeleteSet{fileSet.equalityDeleteSet}
        }
 
-       var (
-               groupKey string
-               groups   map[string][]*equalityDeleteFileSet
-       )
-       groupFiles := make([]*equalityDeleteFileSet, 0, 
len(task.EqualityDeleteFiles))
-
+       files := make([]*equalityDeleteFileSet, 0, 
len(task.EqualityDeleteFiles))
        for _, dataFile := range task.EqualityDeleteFiles {
-               fileSet, ok := perFile[dataFile.FilePath()]
-               if !ok {
-                       continue
-               }
-
-               if groups != nil {
-                       groups[fileSet.groupKey] = 
append(groups[fileSet.groupKey], fileSet)
-               } else if len(groupFiles) == 0 {
-                       groupKey = fileSet.groupKey
-                       groupFiles = append(groupFiles, fileSet)
-               } else if fileSet.groupKey != groupKey {
-                       groups = make(map[string][]*equalityDeleteFileSet, 2)
-                       groups[groupKey] = groupFiles
-                       groupFiles = nil
-                       groups[fileSet.groupKey] = 
append(groups[fileSet.groupKey], fileSet)
-               } else {
-                       groupFiles = append(groupFiles, fileSet)
+               if fileSet, ok := perFile[dataFile.FilePath()]; ok {
+                       files = append(files, fileSet)
                }
        }
 
-       if groups == nil {
-               if len(groupFiles) == 0 {
+       return buildEqualityDeleteSetsForFiles(files, combine)
+}
+
+func buildEqualityDeleteSetsForFiles(
+       files []*equalityDeleteFileSet,
+       combine func([]*equalityDeleteFileSet) *equalityDeleteSet,
+) []*equalityDeleteSet {
+       if len(files) == 0 {
+               return nil
+       }
+       if len(files) == 1 {
+               if len(files[0].keys) == 0 {
                        return nil
                }
 
-               deleteSet := combine(groupFiles)
+               return []*equalityDeleteSet{files[0].equalityDeleteSet}
+       }
+
+       groupKey := files[0].groupKey
+       oneGroup := true
+       for _, file := range files[1:] {
+               if file.groupKey != groupKey {
+                       oneGroup = false
+
+                       break
+               }
+       }
+       if oneGroup {
+               deleteSet := combine(files)
                if len(deleteSet.keys) == 0 {
                        return nil
                }
@@ -677,9 +808,14 @@ func buildEqualityDeleteSetsForTask(
                return []*equalityDeleteSet{deleteSet}
        }
 
+       groups := make(map[string][]*equalityDeleteFileSet, 2)
+       for _, file := range files {
+               groups[file.groupKey] = append(groups[file.groupKey], file)
+       }
+
        sets := make([]*equalityDeleteSet, 0, len(groups))
-       for _, files := range groups {
-               deleteSet := combine(files)
+       for _, groupFiles := range groups {
+               deleteSet := combine(groupFiles)
                if len(deleteSet.keys) > 0 {
                        sets = append(sets, deleteSet)
                }
diff --git a/table/equality_delete_reader_bench_test.go 
b/table/equality_delete_reader_bench_test.go
index 176f27c9f..bcdb536a7 100644
--- a/table/equality_delete_reader_bench_test.go
+++ b/table/equality_delete_reader_bench_test.go
@@ -117,7 +117,44 @@ func BenchmarkReadEqualityDeleteFile(b *testing.B) {
        }
 }
 
-var equalityDeleteLoadingBenchmarkSink int
+var (
+       equalityDeleteLoadingBenchmarkSink  int
+       equalityDeleteMetadataBenchmarkSink int
+)
+
+func BenchmarkLazyEqualityDeleteMetadataSetup(b *testing.B) {
+       tableSchema := iceberg.NewSchema(0,
+               iceberg.NestedField{ID: 1, Name: "id", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+       )
+
+       const taskCount = 10_000
+       for _, uniqueFiles := range []int{1, 100, 1_000} {
+               deleteFiles := make([]iceberg.DataFile, uniqueFiles)
+               for i := range deleteFiles {
+                       deleteFiles[i] = newEqualityDeleteSetAssemblyTestFile(
+                               b, 
fmt.Sprintf("mem://metadata-benchmark/delete-%04d.parquet", i), []int{1})
+               }
+
+               tasks := make([]FileScanTask, taskCount)
+               for i := range tasks {
+                       tasks[i] = FileScanTask{
+                               EqualityDeleteFiles: 
[]iceberg.DataFile{deleteFiles[i%uniqueFiles]},
+                       }
+               }
+
+               b.Run(fmt.Sprintf("tasks=%d/unique_files=%d", taskCount, 
uniqueFiles), func(b *testing.B) {
+                       b.ReportAllocs()
+                       b.ResetTimer()
+                       for b.Loop() {
+                               loader, err := newLazyEqualityDeleteLoader(nil, 
tableSchema, nil, nil, tasks)
+                               if err != nil {
+                                       b.Fatal(err)
+                               }
+                               equalityDeleteMetadataBenchmarkSink = 
len(loader.files)
+                       }
+               })
+       }
+}
 
 type equalityDeleteLoadingBenchmarkInput struct {
        fs          *countingEqualityDeleteOpenFS
@@ -612,3 +649,42 @@ func benchmarkNestedEqualityDeleteSchema(width, depth int) 
(*iceberg.Schema, int
 
        return iceberg.NewSchema(0, fields...), lastLeafID
 }
+
+type equalityMetadataPublicBenchmarkFile struct{ iceberg.DataFile }
+
+func BenchmarkLazyEqualityDeleteMetadataShapes(b *testing.B) {
+       const taskCount = 10000
+       tableSchema := iceberg.NewSchema(0, iceberg.NestedField{ID: 1, Name: 
"id", Type: iceberg.PrimitiveTypes.Int64, Required: true})
+       for _, uniqueFiles := range []int{1, 100, 1000, 10000} {
+               for _, shape := range []string{"shared", "distinct", "public"} {
+                       b.Run(fmt.Sprintf("unique=%d/%s", uniqueFiles, shape), 
func(b *testing.B) {
+                               files := make([]iceberg.DataFile, uniqueFiles)
+                               for i := range files {
+                                       files[i] = 
newEqualityDeleteSetAssemblyTestFile(b, 
fmt.Sprintf("mem://metadata-shapes/delete-%d.parquet", i), []int{1})
+                               }
+                               tasks := make([]FileScanTask, taskCount)
+                               for i := range tasks {
+                                       file := files[i%uniqueFiles]
+                                       if shape != "shared" {
+                                               file = 
newEqualityDeleteSetAssemblyTestFile(b, file.FilePath(), []int{1})
+                                       }
+                                       if shape == "public" {
+                                               file = 
equalityMetadataPublicBenchmarkFile{DataFile: file}
+                                       }
+                                       tasks[i].EqualityDeleteFiles = 
[]iceberg.DataFile{file}
+                               }
+                               b.ReportAllocs()
+                               for b.Loop() {
+                                       loader, err := 
newLazyEqualityDeleteLoader(nil, tableSchema, nil, nil, tasks)
+                                       if err != nil {
+                                               b.Fatal(err)
+                                       }
+                                       if len(loader.files) != uniqueFiles {
+                                               b.Fatalf("got %d unique files, 
want %d", len(loader.files), uniqueFiles)
+                                       }
+                                       equalityDeleteMetadataBenchmarkSink = 
len(loader.files)
+                               }
+                       })
+               }
+       }
+}
diff --git a/table/equality_delete_reader_internal_test.go 
b/table/equality_delete_reader_internal_test.go
index be0b345c0..0368e7fd8 100644
--- a/table/equality_delete_reader_internal_test.go
+++ b/table/equality_delete_reader_internal_test.go
@@ -49,6 +49,26 @@ type countingEqualityDeleteOpenFS struct {
        opens    atomic.Int64
 }
 
+type countingEqualityFieldDataFile struct {
+       iceberg.DataFile
+       equalityFieldIDsCalls         int
+       borrowedEqualityFieldIDsCalls int
+}
+
+func (f *countingEqualityFieldDataFile) EqualityFieldIDs() []int {
+       f.equalityFieldIDsCalls++
+
+       return f.DataFile.EqualityFieldIDs()
+}
+
+func (f *countingEqualityFieldDataFile) DataFileCollectionsRef(_ 
iceinternal.DataFileRef) (
+       map[int]int64, []byte, []int64, []int,
+) {
+       f.borrowedEqualityFieldIDsCalls++
+
+       return iceinternal.BorrowedDataFileCollections(f.DataFile)
+}
+
 func (f *countingEqualityDeleteOpenFS) Open(name string) (iceio.File, error) {
        f.attempts.Add(1)
        file, err := f.MemFS.Open(name)
@@ -265,6 +285,79 @@ func 
TestReadAllEqualityDeleteFilesRejectsEmptyEqualityFieldIDs(t *testing.T) {
        require.ErrorContains(t, err, "empty-equality-fields.parquet")
 }
 
+func TestEqualityDeleteMetadataIsReadOncePerPath(t *testing.T) {
+       t.Parallel()
+
+       schema := iceberg.NewSchema(0,
+               iceberg.NestedField{ID: 1, Name: "id", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+       )
+
+       base := newEqualityDeleteSetAssemblyTestFile(t, 
"mem://metadata-dedup/delete.parquet", []int{1})
+       deleteFile := &countingEqualityFieldDataFile{DataFile: base}
+       tasks := make([]FileScanTask, 100)
+       for i := range tasks {
+               tasks[i] = FileScanTask{EqualityDeleteFiles: 
[]iceberg.DataFile{deleteFile}}
+       }
+
+       loader, err := newLazyEqualityDeleteLoader(iceio.NewMemFS(), schema, 
nil, nil, tasks)
+       require.NoError(t, err)
+       assert.Len(t, loader.files, 1)
+       assert.Equal(t, 1, deleteFile.borrowedEqualityFieldIDsCalls)
+       assert.Zero(t, deleteFile.equalityFieldIDsCalls)
+
+       fs := iceio.NewMemFS()
+       path := "mem://metadata-dedup/eager-delete.parquet"
+       writeEqualityDeleteParquetToMemFS(t, fs, path, `[{"id": 1}]`)
+       base = newEqualityDeleteSetAssemblyTestFile(t, path, []int{1})
+       deleteFile = &countingEqualityFieldDataFile{DataFile: base}
+       tasks = make([]FileScanTask, 100)
+       for i := range tasks {
+               tasks[i] = FileScanTask{EqualityDeleteFiles: 
[]iceberg.DataFile{deleteFile}}
+       }
+
+       perTask, err := readAllEqualityDeleteFiles(t.Context(), fs, schema, 
nil, tasks, 1)
+       require.NoError(t, err)
+       assert.Len(t, perTask, len(tasks))
+       assert.Equal(t, 1, deleteFile.borrowedEqualityFieldIDsCalls)
+       assert.Zero(t, deleteFile.equalityFieldIDsCalls)
+}
+
+func TestLazyEqualityDeleteLoaderNeedsSchemaHistory(t *testing.T) {
+       t.Parallel()
+
+       currentSchema := iceberg.NewSchema(1,
+               iceberg.NestedField{ID: 2, Name: "current", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+       )
+       historicalSchema := iceberg.NewSchema(0,
+               iceberg.NestedField{ID: 1, Name: "id", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+       )
+       fs := iceio.NewMemFS()
+       deletePath := "mem://schema-history/delete.parquet"
+       writeEqualityDeleteParquetToMemFS(t, fs, deletePath, `[{"id": 1}]`)
+       deleteFile := newEqualityDeleteSetAssemblyTestFile(t, deletePath, 
[]int{1})
+       task := FileScanTask{EqualityDeleteFiles: 
[]iceberg.DataFile{deleteFile}}
+
+       loader, err := newLazyEqualityDeleteLoader(fs, currentSchema, nil, nil, 
[]FileScanTask{task})
+       require.NoError(t, err)
+       require.True(t, loader.needsSchemaHistory())
+
+       loader.tableSchemas = []*iceberg.Schema{historicalSchema, currentSchema}
+       sets, err := loader.load(t.Context(), task)
+       require.NoError(t, err)
+       require.Len(t, sets, 1)
+       assert.Equal(t, []int{1}, sets[0].fieldIDs)
+       assert.Equal(t, []string{"id"}, sets[0].colNames)
+       assert.Len(t, sets[0].keys, 1)
+
+       currentDelete := newEqualityDeleteSetAssemblyTestFile(
+               t, "mem://schema-history/current.parquet", []int{2})
+       loader, err = newLazyEqualityDeleteLoader(fs, currentSchema, nil, nil, 
[]FileScanTask{{
+               EqualityDeleteFiles: []iceberg.DataFile{currentDelete},
+       }})
+       require.NoError(t, err)
+       assert.False(t, loader.needsSchemaHistory())
+}
+
 func TestLazyEqualityDeleteLoaderLoadsFilesOnDemand(t *testing.T) {
        t.Parallel()
 

Reply via email to