From abc35da40eff3135267b44b8c9de0fdc22d2218b Mon Sep 17 00:00:00 2001 From: Matt Faltyn Date: Thu, 3 Sep 2026 10:46:43 +0200 Subject: [PATCH 1/5] fix(table): normalize stale last partition ID Signed-off-by: Matt Faltyn --- table/metadata.go | 34 ++++++++++++++++++-------------- table/metadata_preflight_test.go | 10 ++++++++++ 2 files changed, 29 insertions(+), 15 deletions(-) diff --git a/table/metadata.go b/table/metadata.go index c6c5c26d0..fc0d88167 100644 --- a/table/metadata.go +++ b/table/metadata.go @@ -2016,11 +2016,13 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js } lastAssignedID := iceberg.PartitionDataIDStart - 1 + lastPartitionID := lastAssignedID lastPartitionIDSet := false if rawLastPartitionID, ok := metadata["last-partition-id"]; ok { - var lastPartitionID *int - if err := json.Unmarshal(rawLastPartitionID, &lastPartitionID); err == nil && lastPartitionID != nil { - lastAssignedID = max(lastAssignedID, *lastPartitionID) + var parsedLastPartitionID *int + if err := json.Unmarshal(rawLastPartitionID, &parsedLastPartitionID); err == nil && parsedLastPartitionID != nil { + lastPartitionID = *parsedLastPartitionID + lastAssignedID = max(lastAssignedID, lastPartitionID) lastPartitionIDSet = true } } @@ -2042,7 +2044,7 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js } } - if len(missingFields) == 0 { + if len(missingFields) == 0 && (!lastPartitionIDSet || lastPartitionID == lastAssignedID) { return b, nil } @@ -2055,18 +2057,20 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js field["field-id"] = rawFieldID } - if usesSpecList { - rawSpecs, err := json.Marshal(specs) - if err != nil { - return nil, err - } - metadata["partition-specs"] = rawSpecs - } else { - rawFields, err := json.Marshal(specs[0].Fields) - if err != nil { - return nil, err + if len(missingFields) > 0 { + if usesSpecList { + rawSpecs, err := json.Marshal(specs) + if err != nil { + return nil, err + } + metadata["partition-specs"] = rawSpecs + } else { + rawFields, err := json.Marshal(specs[0].Fields) + if err != nil { + return nil, err + } + metadata["partition-spec"] = rawFields } - metadata["partition-spec"] = rawFields } if lastPartitionIDSet { diff --git a/table/metadata_preflight_test.go b/table/metadata_preflight_test.go index f2c48cb0a..0ee2e4c85 100644 --- a/table/metadata_preflight_test.go +++ b/table/metadata_preflight_test.go @@ -85,6 +85,16 @@ func TestParseMetadataBytesAssignsMissingPartitionFieldIDs(t *testing.T) { } } +func TestParseMetadataBytesNormalizesStaleLastPartitionID(t *testing.T) { + data := strings.Replace(ExampleTableMetadataV2, + `"last-partition-id": 1000`, `"last-partition-id": 999`, 1) + + parsed, err := ParseMetadataBytes([]byte(data)) + require.NoError(t, err) + require.NotNil(t, parsed.LastPartitionSpecID()) + assert.Equal(t, 1000, *parsed.LastPartitionSpecID()) +} + func TestParseMetadataBytesRejectsCaseFoldedFormatVersionCollision(t *testing.T) { data := strings.Replace( ExampleTableMetadataV2, From ebf0a161b6fa05cb90f6ff70d1a4fe4b7873c838 Mon Sep 17 00:00:00 2001 From: Matt Faltyn Date: Thu, 3 Sep 2026 20:27:27 +0200 Subject: [PATCH 2/5] fix(table): preserve unassigned partition counters Signed-off-by: Matt Faltyn --- table/metadata.go | 10 +++++-- table/metadata_preflight_test.go | 50 ++++++++++++++++++++++++++++++++ 2 files changed, 57 insertions(+), 3 deletions(-) diff --git a/table/metadata.go b/table/metadata.go index fc0d88167..ae5fd0c77 100644 --- a/table/metadata.go +++ b/table/metadata.go @@ -2016,12 +2016,14 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js } lastAssignedID := iceberg.PartitionDataIDStart - 1 - lastPartitionID := lastAssignedID + lastPartitionID := 0 + normalizedLastPartitionID := 0 lastPartitionIDSet := false if rawLastPartitionID, ok := metadata["last-partition-id"]; ok { var parsedLastPartitionID *int if err := json.Unmarshal(rawLastPartitionID, &parsedLastPartitionID); err == nil && parsedLastPartitionID != nil { lastPartitionID = *parsedLastPartitionID + normalizedLastPartitionID = lastPartitionID lastAssignedID = max(lastAssignedID, lastPartitionID) lastPartitionIDSet = true } @@ -2040,11 +2042,12 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js var fieldID *int if err := json.Unmarshal(rawFieldID, &fieldID); err == nil && fieldID != nil { lastAssignedID = max(lastAssignedID, *fieldID) + normalizedLastPartitionID = max(normalizedLastPartitionID, *fieldID) } } } - if len(missingFields) == 0 && (!lastPartitionIDSet || lastPartitionID == lastAssignedID) { + if len(missingFields) == 0 && (!lastPartitionIDSet || lastPartitionID == normalizedLastPartitionID) { return b, nil } @@ -2055,6 +2058,7 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js return nil, err } field["field-id"] = rawFieldID + normalizedLastPartitionID = max(normalizedLastPartitionID, lastAssignedID) } if len(missingFields) > 0 { @@ -2074,7 +2078,7 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js } if lastPartitionIDSet { - rawLastPartitionID, err := json.Marshal(lastAssignedID) + rawLastPartitionID, err := json.Marshal(normalizedLastPartitionID) if err != nil { return nil, err } diff --git a/table/metadata_preflight_test.go b/table/metadata_preflight_test.go index 0ee2e4c85..c068c71c0 100644 --- a/table/metadata_preflight_test.go +++ b/table/metadata_preflight_test.go @@ -95,6 +95,56 @@ func TestParseMetadataBytesNormalizesStaleLastPartitionID(t *testing.T) { assert.Equal(t, 1000, *parsed.LastPartitionSpecID()) } +func TestAssignMissingPartitionFieldIDsPreservesConsistentMetadata(t *testing.T) { + for _, tt := range []struct { + name string + input string + }{ + { + name: "counter below assignment floor with no fields", + input: `{"last-updated-ms":0,"last-partition-id":0,"partition-specs":[{"spec-id":0,"fields":[]}]}`, + }, + { + name: "counter above greatest field ID", + input: `{"last-updated-ms":0,"last-partition-id":1001,"partition-specs":[{"spec-id":0,"fields":[{"field-id":1000}]}]}`, + }, + } { + t.Run(tt.name, func(t *testing.T) { + normalized, err := assignMissingPartitionFieldIDs([]byte(tt.input)) + require.NoError(t, err) + assert.Equal(t, tt.input, string(normalized)) + }) + } +} + +func TestAssignMissingPartitionFieldIDsNormalizesStaleCounter(t *testing.T) { + input := []byte(`{ + "last-updated-ms": 0, + "last-partition-id": 999, + "partition-specs": [{ + "spec-id": 0, + "fields": [{"field-id": 1000}, {}] + }] + }`) + + normalized, err := assignMissingPartitionFieldIDs(input) + require.NoError(t, err) + + var parsed struct { + LastPartitionID int `json:"last-partition-id"` + Specs []struct { + Fields []struct { + FieldID int `json:"field-id"` + } `json:"fields"` + } `json:"partition-specs"` + } + require.NoError(t, json.Unmarshal(normalized, &parsed)) + assert.Equal(t, 1001, parsed.LastPartitionID) + require.Len(t, parsed.Specs, 1) + require.Len(t, parsed.Specs[0].Fields, 2) + assert.Equal(t, 1001, parsed.Specs[0].Fields[1].FieldID) +} + func TestParseMetadataBytesRejectsCaseFoldedFormatVersionCollision(t *testing.T) { data := strings.Replace( ExampleTableMetadataV2, From e01dbb262a6fd1638fef38c05fe2dde637ac8e83 Mon Sep 17 00:00:00 2001 From: Matt Faltyn Date: Sun, 6 Sep 2026 09:49:05 +0200 Subject: [PATCH 3/5] test(table): cover partition counter normalization Signed-off-by: Matt Faltyn --- table/metadata_preflight_test.go | 23 +++++++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/table/metadata_preflight_test.go b/table/metadata_preflight_test.go index c068c71c0..e5a087d26 100644 --- a/table/metadata_preflight_test.go +++ b/table/metadata_preflight_test.go @@ -22,6 +22,7 @@ import ( "strings" "testing" + "github.com/apache/iceberg-go" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -93,6 +94,15 @@ func TestParseMetadataBytesNormalizesStaleLastPartitionID(t *testing.T) { require.NoError(t, err) require.NotNil(t, parsed.LastPartitionSpecID()) assert.Equal(t, 1000, *parsed.LastPartitionSpecID()) + + update := NewUpdateSpec(New(nil, parsed, "", nil, nil).NewTransaction(), false). + AddField("x", iceberg.BucketTransform{NumBuckets: 16}, "x_bucket") + _, _, err = update.BuildUpdates() + require.NoError(t, err) + updated, err := update.Apply() + require.NoError(t, err) + require.Equal(t, 2, updated.NumFields()) + assert.Equal(t, 1001, updated.Field(1).FieldID) } func TestAssignMissingPartitionFieldIDsPreservesConsistentMetadata(t *testing.T) { @@ -145,6 +155,19 @@ func TestAssignMissingPartitionFieldIDsNormalizesStaleCounter(t *testing.T) { assert.Equal(t, 1001, parsed.Specs[0].Fields[1].FieldID) } +func TestAssignMissingPartitionFieldIDsNormalizesLegacyStaleCounter(t *testing.T) { + input := []byte(`{"last-updated-ms":0,"last-partition-id":8,"partition-specs":[{"spec-id":0,"fields":[{"field-id":9}]}]}`) + + normalized, err := assignMissingPartitionFieldIDs(input) + require.NoError(t, err) + + var parsed struct { + LastPartitionID int `json:"last-partition-id"` + } + require.NoError(t, json.Unmarshal(normalized, &parsed)) + assert.Equal(t, 9, parsed.LastPartitionID) +} + func TestParseMetadataBytesRejectsCaseFoldedFormatVersionCollision(t *testing.T) { data := strings.Replace( ExampleTableMetadataV2, From c9bae4137add146765917fbf9eaa98846cd1fd42 Mon Sep 17 00:00:00 2001 From: Matt Faltyn Date: Wed, 9 Sep 2026 22:16:10 +0200 Subject: [PATCH 4/5] fix(table): allocate partition IDs from history Signed-off-by: Matt Faltyn --- table/metadata.go | 61 +++++++++++++------------ table/metadata_builder_internal_test.go | 27 +++++++++++ table/metadata_preflight_test.go | 59 ++++++------------------ table/update_spec.go | 8 +--- 4 files changed, 74 insertions(+), 81 deletions(-) diff --git a/table/metadata.go b/table/metadata.go index ae5fd0c77..d5554e5e2 100644 --- a/table/metadata.go +++ b/table/metadata.go @@ -252,9 +252,9 @@ type Metadata interface { // DefaultPartitionSpec is the ID of the current spec that writerFactory should // use by default. DefaultPartitionSpec() int - // LastPartitionSpecID is the highest assigned partition field ID across - // all partition specs for the table. This is used to ensure partition - // fields are always assigned an unused ID when evolving specs. + // LastPartitionSpecID returns the persisted last assigned partition field ID. + // Allocation also scans partition spec history because metadata written by + // another client may contain a stale counter. LastPartitionSpecID() *int // Snapshots returns the list of valid snapshots. Valid snapshots are // snapshots for which all data files exist in the file system. A data @@ -700,6 +700,18 @@ func (b *MetadataBuilder) AddSchema(schema *iceberg.Schema) error { return nil } +func partitionFieldIDFloor(lastPartitionID *int, specs []iceberg.PartitionSpec) int { + floor := partitionFieldStartID - 1 + if lastPartitionID != nil { + floor = max(floor, *lastPartitionID) + } + for _, spec := range specs { + floor = max(floor, spec.LastAssignedFieldID()) + } + + return floor +} + func (b *MetadataBuilder) AddPartitionSpec(spec *iceberg.PartitionSpec, initial bool) error { newSpecID := b.reuseOrCreateNewPartitionSpecID(*spec) curSchema := b.CurrentSchema() @@ -707,7 +719,8 @@ func (b *MetadataBuilder) AddPartitionSpec(spec *iceberg.PartitionSpec, initial return errors.New("can't add sort order with no current schema") } - freshSpec, err := spec.BindToSchema(curSchema, b.lastPartitionID, &newSpecID) + fieldIDFloor := partitionFieldIDFloor(b.lastPartitionID, b.specs) + freshSpec, err := spec.BindToSchema(curSchema, &fieldIDFloor, &newSpecID) if err != nil { return err } @@ -2016,15 +2029,11 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js } lastAssignedID := iceberg.PartitionDataIDStart - 1 - lastPartitionID := 0 - normalizedLastPartitionID := 0 lastPartitionIDSet := false if rawLastPartitionID, ok := metadata["last-partition-id"]; ok { - var parsedLastPartitionID *int - if err := json.Unmarshal(rawLastPartitionID, &parsedLastPartitionID); err == nil && parsedLastPartitionID != nil { - lastPartitionID = *parsedLastPartitionID - normalizedLastPartitionID = lastPartitionID - lastAssignedID = max(lastAssignedID, lastPartitionID) + var lastPartitionID *int + if err := json.Unmarshal(rawLastPartitionID, &lastPartitionID); err == nil && lastPartitionID != nil { + lastAssignedID = max(lastAssignedID, *lastPartitionID) lastPartitionIDSet = true } } @@ -2042,12 +2051,11 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js var fieldID *int if err := json.Unmarshal(rawFieldID, &fieldID); err == nil && fieldID != nil { lastAssignedID = max(lastAssignedID, *fieldID) - normalizedLastPartitionID = max(normalizedLastPartitionID, *fieldID) } } } - if len(missingFields) == 0 && (!lastPartitionIDSet || lastPartitionID == normalizedLastPartitionID) { + if len(missingFields) == 0 { return b, nil } @@ -2058,27 +2066,24 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js return nil, err } field["field-id"] = rawFieldID - normalizedLastPartitionID = max(normalizedLastPartitionID, lastAssignedID) } - if len(missingFields) > 0 { - if usesSpecList { - rawSpecs, err := json.Marshal(specs) - if err != nil { - return nil, err - } - metadata["partition-specs"] = rawSpecs - } else { - rawFields, err := json.Marshal(specs[0].Fields) - if err != nil { - return nil, err - } - metadata["partition-spec"] = rawFields + if usesSpecList { + rawSpecs, err := json.Marshal(specs) + if err != nil { + return nil, err + } + metadata["partition-specs"] = rawSpecs + } else { + rawFields, err := json.Marshal(specs[0].Fields) + if err != nil { + return nil, err } + metadata["partition-spec"] = rawFields } if lastPartitionIDSet { - rawLastPartitionID, err := json.Marshal(normalizedLastPartitionID) + rawLastPartitionID, err := json.Marshal(lastAssignedID) if err != nil { return nil, err } diff --git a/table/metadata_builder_internal_test.go b/table/metadata_builder_internal_test.go index 9825da0d7..0dba74ab6 100644 --- a/table/metadata_builder_internal_test.go +++ b/table/metadata_builder_internal_test.go @@ -357,6 +357,33 @@ func TestAddRemovePartitionSpec(t *testing.T) { require.ErrorContains(t, err, "id 1") } +func TestAddPartitionSpecAllocatesAfterHistoricalFieldID(t *testing.T) { + data := strings.Replace(ExampleTableMetadataV2, + `"last-partition-id": 1000`, `"last-partition-id": 999`, 1) + require.Contains(t, data, `"last-partition-id": 999`) + + metadata, err := ParseMetadataBytes([]byte(data)) + require.NoError(t, err) + builder, err := MetadataBuilderFromBase(metadata, "") + require.NoError(t, err) + + addedSpec, err := iceberg.NewPartitionSpecOpts( + iceberg.WithSpecID(1), + iceberg.AddPartitionFieldBySourceID(1, "x_bucket", iceberg.BucketTransform{NumBuckets: 16}, builder.CurrentSchema(), nil), + ) + require.NoError(t, err) + require.NoError(t, builder.AddPartitionSpec(&addedSpec, false)) + + rebuilt, err := builder.Build() + require.NoError(t, err) + added := rebuilt.PartitionSpecByID(1) + require.NotNil(t, added) + require.Equal(t, 1, added.NumFields()) + assert.Equal(t, 1001, added.Field(0).FieldID) + require.NotNil(t, rebuilt.LastPartitionSpecID()) + assert.Equal(t, 1001, *rebuilt.LastPartitionSpecID()) +} + func TestRemovePartitionSpecsNoMatchDoesNotUpdate(t *testing.T) { for _, count := range []int{1, 2, 8, 9, 16, 17, 32, 33, 64} { t.Run(fmt.Sprintf("requested=%d", count), func(t *testing.T) { diff --git a/table/metadata_preflight_test.go b/table/metadata_preflight_test.go index e5a087d26..22a2cff12 100644 --- a/table/metadata_preflight_test.go +++ b/table/metadata_preflight_test.go @@ -86,23 +86,31 @@ func TestParseMetadataBytesAssignsMissingPartitionFieldIDs(t *testing.T) { } } -func TestParseMetadataBytesNormalizesStaleLastPartitionID(t *testing.T) { +func TestParseMetadataBytesPreservesStaleLastPartitionIDForCommit(t *testing.T) { data := strings.Replace(ExampleTableMetadataV2, `"last-partition-id": 1000`, `"last-partition-id": 999`, 1) + data = strings.Replace(data, + `"default-spec-id": 0,`, `"default-spec-id": 1,`, 1) + data = strings.Replace(data, + `"partition-specs": [{"spec-id": 0, "fields": [{"name": "x", "transform": "identity", "source-id": 1, "field-id": 1000}]}],`, + `"partition-specs": [{"spec-id": 0, "fields": [{"name": "x", "transform": "identity", "source-id": 1, "field-id": 1000}]}, {"spec-id": 1, "fields": []}],`, 1) + require.Contains(t, data, `"last-partition-id": 999`) + require.Contains(t, data, `"spec-id": 1`) parsed, err := ParseMetadataBytes([]byte(data)) require.NoError(t, err) require.NotNil(t, parsed.LastPartitionSpecID()) - assert.Equal(t, 1000, *parsed.LastPartitionSpecID()) + assert.Equal(t, 999, *parsed.LastPartitionSpecID()) update := NewUpdateSpec(New(nil, parsed, "", nil, nil).NewTransaction(), false). AddField("x", iceberg.BucketTransform{NumBuckets: 16}, "x_bucket") - _, _, err = update.BuildUpdates() + _, requirements, err := update.BuildUpdates() require.NoError(t, err) + assert.Equal(t, []int{999}, lastAssignedPartitionAssertions(requirements)) updated, err := update.Apply() require.NoError(t, err) - require.Equal(t, 2, updated.NumFields()) - assert.Equal(t, 1001, updated.Field(1).FieldID) + require.Equal(t, 1, updated.NumFields()) + assert.Equal(t, 1001, updated.Field(0).FieldID) } func TestAssignMissingPartitionFieldIDsPreservesConsistentMetadata(t *testing.T) { @@ -127,47 +135,6 @@ func TestAssignMissingPartitionFieldIDsPreservesConsistentMetadata(t *testing.T) } } -func TestAssignMissingPartitionFieldIDsNormalizesStaleCounter(t *testing.T) { - input := []byte(`{ - "last-updated-ms": 0, - "last-partition-id": 999, - "partition-specs": [{ - "spec-id": 0, - "fields": [{"field-id": 1000}, {}] - }] - }`) - - normalized, err := assignMissingPartitionFieldIDs(input) - require.NoError(t, err) - - var parsed struct { - LastPartitionID int `json:"last-partition-id"` - Specs []struct { - Fields []struct { - FieldID int `json:"field-id"` - } `json:"fields"` - } `json:"partition-specs"` - } - require.NoError(t, json.Unmarshal(normalized, &parsed)) - assert.Equal(t, 1001, parsed.LastPartitionID) - require.Len(t, parsed.Specs, 1) - require.Len(t, parsed.Specs[0].Fields, 2) - assert.Equal(t, 1001, parsed.Specs[0].Fields[1].FieldID) -} - -func TestAssignMissingPartitionFieldIDsNormalizesLegacyStaleCounter(t *testing.T) { - input := []byte(`{"last-updated-ms":0,"last-partition-id":8,"partition-specs":[{"spec-id":0,"fields":[{"field-id":9}]}]}`) - - normalized, err := assignMissingPartitionFieldIDs(input) - require.NoError(t, err) - - var parsed struct { - LastPartitionID int `json:"last-partition-id"` - } - require.NoError(t, json.Unmarshal(normalized, &parsed)) - assert.Equal(t, 9, parsed.LastPartitionID) -} - func TestParseMetadataBytesRejectsCaseFoldedFormatVersionCollision(t *testing.T) { data := strings.Replace( ExampleTableMetadataV2, diff --git a/table/update_spec.go b/table/update_spec.go index a73e2699f..dae04d0e1 100644 --- a/table/update_spec.go +++ b/table/update_spec.go @@ -122,15 +122,9 @@ func NewUpdateSpec(t *Transaction, caseSensitive bool) *UpdateSpec { nameToField[partitionField.Name] = partitionField } us.schema = stagedMeta.CurrentSchema() - lastAssignedFieldId := us.meta.LastPartitionSpecID() - if lastAssignedFieldId == nil { - v := iceberg.PartitionDataIDStart - 1 - lastAssignedFieldId = &v - } - us.nameToField = nameToField us.transformToField = transformToField - us.lastAssignedFieldId = *lastAssignedFieldId + us.lastAssignedFieldId = partitionFieldIDFloor(us.meta.LastPartitionSpecID(), us.meta.PartitionSpecs()) return us } From 96d47a607ba9052e29fb622eaa08b711026ba821 Mon Sep 17 00:00:00 2001 From: Matt Faltyn Date: Fri, 11 Sep 2026 11:32:43 +0200 Subject: [PATCH 5/5] fix(table): repair partition counter write-back Signed-off-by: Matt Faltyn --- table/metadata.go | 15 +++---- table/metadata_builder_internal_test.go | 55 ++++++++++++++++++++----- 2 files changed, 50 insertions(+), 20 deletions(-) diff --git a/table/metadata.go b/table/metadata.go index d5554e5e2..fca5a3087 100644 --- a/table/metadata.go +++ b/table/metadata.go @@ -252,9 +252,8 @@ type Metadata interface { // DefaultPartitionSpec is the ID of the current spec that writerFactory should // use by default. DefaultPartitionSpec() int - // LastPartitionSpecID returns the persisted last assigned partition field ID. - // Allocation also scans partition spec history because metadata written by - // another client may contain a stale counter. + // LastPartitionSpecID returns the persisted last assigned partition field ID, + // which may be stale relative to partition spec history. LastPartitionSpecID() *int // Snapshots returns the list of valid snapshots. Valid snapshots are // snapshots for which all data files exist in the file system. A data @@ -700,8 +699,10 @@ func (b *MetadataBuilder) AddSchema(schema *iceberg.Schema) error { return nil } +// partitionFieldIDFloor returns an allocation floor that accounts for the +// persisted counter and all partition field IDs in spec history. func partitionFieldIDFloor(lastPartitionID *int, specs []iceberg.PartitionSpec) int { - floor := partitionFieldStartID - 1 + floor := iceberg.PartitionDataIDStart - 1 if lastPartitionID != nil { floor = max(floor, *lastPartitionID) } @@ -749,11 +750,7 @@ func (b *MetadataBuilder) AddPartitionSpec(spec *iceberg.PartitionSpec, initial } - prev := partitionFieldStartID - 1 - if b.lastPartitionID != nil { - prev = *b.lastPartitionID - } - lastPartitionID := max(maxFieldID, prev) + lastPartitionID := max(maxFieldID, fieldIDFloor) var specs []iceberg.PartitionSpec if initial { diff --git a/table/metadata_builder_internal_test.go b/table/metadata_builder_internal_test.go index 0dba74ab6..97e0e8c45 100644 --- a/table/metadata_builder_internal_test.go +++ b/table/metadata_builder_internal_test.go @@ -362,26 +362,59 @@ func TestAddPartitionSpecAllocatesAfterHistoricalFieldID(t *testing.T) { `"last-partition-id": 1000`, `"last-partition-id": 999`, 1) require.Contains(t, data, `"last-partition-id": 999`) + metadata, err := ParseMetadataBytes([]byte(data)) + require.NoError(t, err) + + for _, tt := range []struct { + name string + dropCounter bool + }{ + {name: "stale counter"}, + {name: "nil counter", dropCounter: true}, + } { + t.Run(tt.name, func(t *testing.T) { + builder, err := MetadataBuilderFromBase(metadata, "") + require.NoError(t, err) + if tt.dropCounter { + builder.lastPartitionID = nil + } + + addedSpec, err := iceberg.NewPartitionSpecOpts( + iceberg.WithSpecID(1), + iceberg.AddPartitionFieldBySourceID(1, "x_bucket", iceberg.BucketTransform{NumBuckets: 16}, builder.CurrentSchema(), nil), + ) + require.NoError(t, err) + require.NoError(t, builder.AddPartitionSpec(&addedSpec, false)) + + rebuilt, err := builder.Build() + require.NoError(t, err) + added := rebuilt.PartitionSpecByID(1) + require.NotNil(t, added) + require.Equal(t, 1, added.NumFields()) + assert.Equal(t, 1001, added.Field(0).FieldID) + require.NotNil(t, rebuilt.LastPartitionSpecID()) + assert.Equal(t, 1001, *rebuilt.LastPartitionSpecID()) + }) + } +} + +func TestAddEmptyPartitionSpecRepairsStaleLastPartitionID(t *testing.T) { + data := strings.Replace(ExampleTableMetadataV2, + `"last-partition-id": 1000`, `"last-partition-id": 999`, 1) + require.Contains(t, data, `"last-partition-id": 999`) + metadata, err := ParseMetadataBytes([]byte(data)) require.NoError(t, err) builder, err := MetadataBuilderFromBase(metadata, "") require.NoError(t, err) - addedSpec, err := iceberg.NewPartitionSpecOpts( - iceberg.WithSpecID(1), - iceberg.AddPartitionFieldBySourceID(1, "x_bucket", iceberg.BucketTransform{NumBuckets: 16}, builder.CurrentSchema(), nil), - ) - require.NoError(t, err) - require.NoError(t, builder.AddPartitionSpec(&addedSpec, false)) + emptySpec := iceberg.NewPartitionSpecID(1) + require.NoError(t, builder.AddPartitionSpec(&emptySpec, false)) rebuilt, err := builder.Build() require.NoError(t, err) - added := rebuilt.PartitionSpecByID(1) - require.NotNil(t, added) - require.Equal(t, 1, added.NumFields()) - assert.Equal(t, 1001, added.Field(0).FieldID) require.NotNil(t, rebuilt.LastPartitionSpecID()) - assert.Equal(t, 1001, *rebuilt.LastPartitionSpecID()) + assert.Equal(t, 1000, *rebuilt.LastPartitionSpecID()) } func TestRemovePartitionSpecsNoMatchDoesNotUpdate(t *testing.T) {