diff --git a/table/metadata.go b/table/metadata.go index c6c5c26d0..ae5fd0c77 100644 --- a/table/metadata.go +++ b/table/metadata.go @@ -2016,11 +2016,15 @@ 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 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 + normalizedLastPartitionID = lastPartitionID + lastAssignedID = max(lastAssignedID, lastPartitionID) lastPartitionIDSet = true } } @@ -2038,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 { + if len(missingFields) == 0 && (!lastPartitionIDSet || lastPartitionID == normalizedLastPartitionID) { return b, nil } @@ -2053,24 +2058,27 @@ func assignMissingPartitionFieldIDsFromMetadata(b []byte, metadata map[string]js return nil, err } field["field-id"] = rawFieldID + normalizedLastPartitionID = max(normalizedLastPartitionID, lastAssignedID) } - 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 { - 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 f2c48cb0a..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" ) @@ -85,6 +86,88 @@ 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()) + + 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) { + 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 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,