Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 0 additions & 51 deletions catalog/catalog.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,12 +32,10 @@ import (
"errors"
"fmt"
"iter"
"maps"
"strings"
"unicode"

"github.com/apache/iceberg-go"
iceinternal "github.com/apache/iceberg-go/internal"
"github.com/apache/iceberg-go/table"
)

Expand Down Expand Up @@ -357,52 +355,3 @@ func WithViewProperties(config iceberg.Properties) CreateViewOpt {
cfg.Properties = config
}
}

//lint:ignore U1000 this is linked to by catalogs via go:linkname but we don't want to export it
func checkForOverlap(removals []string, updates iceberg.Properties) error {
overlap := []string{}
for _, key := range removals {
if _, ok := updates[key]; ok {
overlap = append(overlap, key)
}
}
if len(overlap) > 0 {
return fmt.Errorf("conflict between removals and updates for keys: %v", overlap)
}

return nil
}

//lint:ignore U1000 this is linked to by catalogs via go:linkname but we don't want to export it
func getUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, PropertiesUpdateSummary, error) {
if err := checkForOverlap(removals, updates); err != nil {
return nil, PropertiesUpdateSummary{}, err
}
var (
updatedProps = maps.Clone(currentProps)
removed = make([]string, 0, len(removals))
updated = make([]string, 0, len(updates))
)

for _, key := range removals {
if _, exists := updatedProps[key]; exists {
delete(updatedProps, key)
removed = append(removed, key)
}
}

for key, value := range updates {
if updatedProps[key] != value {
updated = append(updated, key)
updatedProps[key] = value
}
}

summary := PropertiesUpdateSummary{
Removed: removed,
Updated: updated,
Missing: iceinternal.Difference(removals, removed),
}

return updatedProps, summary, nil
}
9 changes: 1 addition & 8 deletions catalog/glue/glue.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@ import (
"strconv"
"strings"
"time"
_ "unsafe"

"github.com/apache/iceberg-go"
"github.com/apache/iceberg-go/catalog"
Expand Down Expand Up @@ -766,12 +765,6 @@ func (c *Catalog) LoadNamespaceProperties(ctx context.Context, namespace table.I
return props, nil
}

// avoid circular dependency while still avoiding having to export the getUpdatedPropsAndUpdateSummary function
// so that we can re-use it in the catalog implementations without duplicating the code.

//go:linkname getUpdatedPropsAndUpdateSummary github.com/apache/iceberg-go/catalog.getUpdatedPropsAndUpdateSummary
func getUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, catalog.PropertiesUpdateSummary, error)

// UpdateNamespaceProperties updates the properties of an Iceberg namespace in the Glue catalog.
// The removals list contains the keys to remove, and the updates map contains the keys and values to update.
func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table.Identifier,
Expand All @@ -782,7 +775,7 @@ func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table
return catalog.PropertiesUpdateSummary{}, err
}

updatedProperties, propertiesUpdateSummary, err := getUpdatedPropsAndUpdateSummary(currentProps, removals, updates)
updatedProperties, propertiesUpdateSummary, err := internal.GetUpdatedPropsAndUpdateSummary(currentProps, removals, updates)
if err != nil {
return catalog.PropertiesUpdateSummary{}, err
}
Expand Down
6 changes: 1 addition & 5 deletions catalog/hive/hive.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,6 @@ import (
"log"
"maps"
"strings"
_ "unsafe"

"github.com/apache/iceberg-go"
"github.com/apache/iceberg-go/catalog"
Expand Down Expand Up @@ -934,9 +933,6 @@ func (c *Catalog) LoadNamespaceProperties(ctx context.Context, namespace table.I
return props, nil
}

//go:linkname getUpdatedPropsAndUpdateSummary github.com/apache/iceberg-go/catalog.getUpdatedPropsAndUpdateSummary
func getUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, catalog.PropertiesUpdateSummary, error)

// UpdateNamespaceProperties updates the properties for a namespace.
func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table.Identifier,
removals []string, updates iceberg.Properties,
Expand All @@ -946,7 +942,7 @@ func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table
return catalog.PropertiesUpdateSummary{}, err
}

updatedProperties, propertiesUpdateSummary, err := getUpdatedPropsAndUpdateSummary(currentProps, removals, updates)
updatedProperties, propertiesUpdateSummary, err := internal.GetUpdatedPropsAndUpdateSummary(currentProps, removals, updates)
if err != nil {
return catalog.PropertiesUpdateSummary{}, err
}
Expand Down
58 changes: 54 additions & 4 deletions catalog/internal/utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ import (

"github.com/apache/iceberg-go"
"github.com/apache/iceberg-go/catalog"
"github.com/apache/iceberg-go/internal"
iceinternal "github.com/apache/iceberg-go/internal"
icebergio "github.com/apache/iceberg-go/io"
"github.com/apache/iceberg-go/table"
"github.com/google/uuid"
Expand All @@ -52,21 +52,21 @@ func WriteTableMetadata(metadata table.Metadata, fs icebergio.WriteFileIO, loc s
if err != nil {
return err
}
defer internal.CheckedClose(out, &err)
defer iceinternal.CheckedClose(out, &err)

var writer io.Writer = out
switch compression {
case table.MetadataCompressionCodecGzip:
gzw := gzip.NewWriter(out)
writer = gzw
defer internal.CheckedClose(gzw, &err)
defer iceinternal.CheckedClose(gzw, &err)
case table.MetadataCompressionCodecZstd:
enc, zErr := zstd.NewWriter(out)
if zErr != nil {
return zErr
}
writer = enc
defer internal.CheckedClose(enc, &err)
defer iceinternal.CheckedClose(enc, &err)
}

err = json.NewEncoder(writer).Encode(metadata)
Expand Down Expand Up @@ -324,3 +324,53 @@ func UpdateAndStageTable(ctx context.Context, catprops iceberg.Properties, curre
),
}, nil
}

func checkForOverlap(removals []string, updates iceberg.Properties) error {
overlap := []string{}
for _, key := range removals {
if _, ok := updates[key]; ok {
overlap = append(overlap, key)
}
}
if len(overlap) > 0 {
return fmt.Errorf("conflict between removals and updates for keys: %v", overlap)
}

return nil
}

// GetUpdatedPropsAndUpdateSummary applies removals and updates to currentProps
// and returns the updated properties alongside a summary of the changes. It is
// shared by the catalog backend implementations.
func GetUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, catalog.PropertiesUpdateSummary, error) {
Comment thread
nssalian marked this conversation as resolved.
if err := checkForOverlap(removals, updates); err != nil {
return nil, catalog.PropertiesUpdateSummary{}, err
}
var (
updatedProps = maps.Clone(currentProps)
removed = make([]string, 0, len(removals))
updated = make([]string, 0, len(updates))
)

for _, key := range removals {
if _, exists := updatedProps[key]; exists {
delete(updatedProps, key)
removed = append(removed, key)
}
}

for key, value := range updates {
if updatedProps[key] != value {
Comment thread
nssalian marked this conversation as resolved.
updated = append(updated, key)
updatedProps[key] = value
}
}

summary := catalog.PropertiesUpdateSummary{
Removed: removed,
Updated: updated,
Missing: iceinternal.Difference(removals, removed),
}

return updatedProps, summary, nil
}
137 changes: 137 additions & 0 deletions catalog/internal/utils_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
// 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 internal

import (
"maps"
"testing"

"github.com/apache/iceberg-go"
"github.com/apache/iceberg-go/catalog"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

func TestGetUpdatedPropsAndUpdateSummary(t *testing.T) {
tests := []struct {
name string
current iceberg.Properties
removals []string
updates iceberg.Properties
wantErrContains string
wantProps iceberg.Properties
wantRemoved []string
wantUpdated []string
wantMissing []string
}{
{
name: "removal and update of the same key conflicts",
current: iceberg.Properties{"a": "1"},
removals: []string{"a"},
updates: iceberg.Properties{"a": "2"},
wantErrContains: "conflict between removals and updates for keys: [a]",
},
{
name: "existing key is removed",
current: iceberg.Properties{"a": "1", "b": "2"},
removals: []string{"a"},
wantProps: iceberg.Properties{"b": "2"},
wantRemoved: []string{"a"},
wantUpdated: []string{},
wantMissing: []string{},
},
{
name: "removal of absent key is reported as missing",
current: iceberg.Properties{"a": "1"},
removals: []string{"ghost"},
wantProps: iceberg.Properties{"a": "1"},
wantRemoved: []string{},
wantUpdated: []string{},
wantMissing: []string{"ghost"},
},
{
name: "existing key gets a new value",
current: iceberg.Properties{"a": "old"},
updates: iceberg.Properties{"a": "new"},
wantProps: iceberg.Properties{"a": "new"},
wantRemoved: []string{},
wantUpdated: []string{"a"},
wantMissing: []string{},
},
{
name: "new key is added",
current: iceberg.Properties{"a": "1"},
updates: iceberg.Properties{"b": "2"},
wantProps: iceberg.Properties{"a": "1", "b": "2"},
wantRemoved: []string{},
wantUpdated: []string{"b"},
wantMissing: []string{},
},
{
name: "update with unchanged value is not reported",
current: iceberg.Properties{"a": "1"},
updates: iceberg.Properties{"a": "1"},
wantProps: iceberg.Properties{"a": "1"},
wantRemoved: []string{},
wantUpdated: []string{},
wantMissing: []string{},
},
{
name: "no removals or updates is a no-op",
current: iceberg.Properties{"a": "1"},
wantProps: iceberg.Properties{"a": "1"},
wantRemoved: []string{},
wantUpdated: []string{},
wantMissing: []string{},
},
{
name: "add, update, remove and missing in one call",
current: iceberg.Properties{"keep": "v", "drop": "v", "upd": "old"},
removals: []string{"drop", "ghost"},
updates: iceberg.Properties{"upd": "new", "add": "v", "keep": "v"},
wantProps: iceberg.Properties{"keep": "v", "upd": "new", "add": "v"},
wantRemoved: []string{"drop"},
wantUpdated: []string{"upd", "add"},
wantMissing: []string{"ghost"},
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
before := maps.Clone(tt.current)

props, summary, err := GetUpdatedPropsAndUpdateSummary(tt.current, tt.removals, tt.updates)

assert.Equal(t, before, tt.current, "currentProps must not be mutated")

if tt.wantErrContains != "" {
require.ErrorContains(t, err, tt.wantErrContains)
assert.Nil(t, props)
assert.Equal(t, catalog.PropertiesUpdateSummary{}, summary)

return
}

require.NoError(t, err)
assert.Equal(t, tt.wantProps, props)
assert.ElementsMatch(t, tt.wantRemoved, summary.Removed)
assert.ElementsMatch(t, tt.wantUpdated, summary.Updated)
assert.ElementsMatch(t, tt.wantMissing, summary.Missing)
})
}
}
9 changes: 1 addition & 8 deletions catalog/sql/sql.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@ import (
"strings"
"sync"
"time"
_ "unsafe"

"github.com/apache/iceberg-go"
"github.com/apache/iceberg-go/catalog"
Expand Down Expand Up @@ -1518,12 +1517,6 @@ func (c *Catalog) ListNamespaces(ctx context.Context, parent table.Identifier) (
return ret, nil
}

// avoid circular dependency while still avoiding having to export the getUpdatedPropsAndUpdateSummary function
// so that we can re-use it in the catalog implementations without duplicating the code.

//go:linkname getUpdatedPropsAndUpdateSummary github.com/apache/iceberg-go/catalog.getUpdatedPropsAndUpdateSummary
func getUpdatedPropsAndUpdateSummary(currentProps iceberg.Properties, removals []string, updates iceberg.Properties) (iceberg.Properties, catalog.PropertiesUpdateSummary, error)

func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table.Identifier, removals []string, updates iceberg.Properties) (catalog.PropertiesUpdateSummary, error) {
var summary catalog.PropertiesUpdateSummary
if err := checkValidNamespace(namespace); err != nil {
Expand Down Expand Up @@ -1555,7 +1548,7 @@ func (c *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table
currentProps[prop.PropertyKey] = prop.PropertyValue.String
}

_, nextSummary, err := getUpdatedPropsAndUpdateSummary(currentProps, removals, updates)
_, nextSummary, err := internal.GetUpdatedPropsAndUpdateSummary(currentProps, removals, updates)
if err != nil {
return err
}
Expand Down
Loading