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