This is an automated email from the ASF dual-hosted git repository.
zeroshade 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 103620f10 fix(table): stop MetadataBuilder.Build from growing the
metadata log (#2095)
103620f10 is described below
commit 103620f105776726514b59f27938fd3155390612
Author: Animesh Kumar <[email protected]>
AuthorDate: Fri Oct 2 23:07:41 2026 +0530
fix(table): stop MetadataBuilder.Build from growing the metadata log (#2095)
* fix(table): stop MetadataBuilder.Build from growing the metadata log
buildCommonMetadata appended the previous metadata file to b.metadataLog
and trimmed it in place, so every Build() on a builder with changes left
one more identical entry behind. Transactions call Build() several times
before committing: a copy-on-write delete builds once per rewritten file,
and StagedTable() and Transaction.Scan build again. A staged table after a
delete that rewrote two files listed the same previous metadata file five
times. Persisted metadata was not affected, because catalogs rebuild from
the base and the updates.
Extend and trim a copy of the log instead, so Build() leaves the builder
unchanged and returns the same log however often it runs.
Closes #2091
Signed-off-by: Animesh Kumar <[email protected]>
* refactor(table): share the metadata log trim and cover the no-changes
build
Address review on #2095: clone the log on both paths of
buildCommonMetadata, move the trim into a pure trimMetadataLog helper that
TrimMetadataLogs also uses, assert on the builder's own log, add a test for a
build without changes, and reword the copy-on-write test comment.
Signed-off-by: Animesh Kumar <[email protected]>
---------
Signed-off-by: Animesh Kumar <[email protected]>
---
table/copy_on_write_metadata_log_test.go | 67 ++++++++++++++++++++++++++++++++
table/metadata.go | 24 ++++++++----
table/metadata_builder_internal_test.go | 54 +++++++++++++++++++++++++
3 files changed, 137 insertions(+), 8 deletions(-)
diff --git a/table/copy_on_write_metadata_log_test.go
b/table/copy_on_write_metadata_log_test.go
new file mode 100644
index 000000000..0858e4235
--- /dev/null
+++ b/table/copy_on_write_metadata_log_test.go
@@ -0,0 +1,67 @@
+// 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_test
+
+import (
+ "context"
+ "path/filepath"
+ "slices"
+ "testing"
+
+ "github.com/apache/iceberg-go"
+ iceio "github.com/apache/iceberg-go/io"
+ "github.com/apache/iceberg-go/table"
+ "github.com/stretchr/testify/require"
+)
+
+// A transaction builds its metadata several times before it commits (for a
+// copy-on-write delete: while classifying files, once per rewritten file, and
again
+// for StagedTable). However many builds there are, the staged table must
carry a
+// single previous-metadata entry.
+func TestCopyOnWriteDeleteStagesSinglePreviousMetadataLogEntry(t *testing.T) {
+ ctx := context.Background()
+ location := filepath.ToSlash(t.TempDir())
+ 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.String, Required: false},
+ )
+ meta, err := table.NewMetadata(schema, iceberg.UnpartitionedSpec,
table.UnsortedSortOrder, location,
+ iceberg.Properties{
+ table.PropertyFormatVersion: "2",
+ table.WriteDeleteModeKey: table.WriteModeCopyOnWrite,
+ })
+ require.NoError(t, err)
+
+ metaLoc := location + "/metadata/v1.metadata.json"
+ fsF := func(context.Context) (iceio.IO, error) { return
iceio.LocalFS{}, nil }
+ cat := &concurrentTestCatalog{metadata: meta, location: metaLoc, fsF:
fsF}
+ tbl := table.New(table.Identifier{"db", "cow_metadata_log"}, meta,
metaLoc, fsF, cat)
+ tbl = appendTenRows(t, appendTenRows(t, tbl))
+
+ tasks, err := tbl.Scan().PlanFiles(ctx)
+ require.NoError(t, err)
+ require.Len(t, tasks, 2, "setup: the delete must rewrite more than one
file")
+
+ txn := tbl.NewTransaction()
+ require.NoError(t, txn.Delete(ctx,
iceberg.EqualTo(iceberg.Reference("id"), int64(2)), nil))
+
+ staged, err := txn.StagedTable()
+ require.NoError(t, err)
+ require.Len(t, slices.Collect(staged.Metadata().PreviousFiles()), 1,
+ "each Build must not add another previous-metadata entry")
+}
diff --git a/table/metadata.go b/table/metadata.go
index 2ad17502f..bf1dc9e26 100644
--- a/table/metadata.go
+++ b/table/metadata.go
@@ -1517,12 +1517,14 @@ func (b *MetadataBuilder) buildCommonMetadata()
(*commonMetadata, error) {
b.lastUpdatedMS = time.Now().UnixMilli()
}
+ // Build may run more than once on the same builder, so the log is
extended on a
+ // copy. Appending to b.metadataLog would add the previous file again
on every call.
+ metadataLog := slices.Clone(b.metadataLog)
if b.previousFileEntry != nil && b.HasChanges() {
maxMetadataLogEntries := max(1,
b.props.GetInt(
MetadataPreviousVersionsMaxKey,
MetadataPreviousVersionsMaxDefault))
- b.AppendMetadataLog(*b.previousFileEntry)
- b.TrimMetadataLogs(maxMetadataLogEntries)
+ metadataLog = trimMetadataLog(append(metadataLog,
*b.previousFileEntry), maxMetadataLogEntries)
}
return &commonMetadata{
@@ -1543,7 +1545,7 @@ func (b *MetadataBuilder) buildCommonMetadata()
(*commonMetadata, error) {
snapshotIndex: b.snapshotIndex,
CurrentSnapshotID: b.currentSnapshotID,
SnapshotLog: b.snapshotLog,
- MetadataLog: b.metadataLog,
+ MetadataLog: metadataLog,
SortOrderList: b.sortOrderList,
DefaultSortOrderID: b.defaultSortOrderID,
SnapshotRefs: b.refs,
@@ -1663,15 +1665,21 @@ func (b *MetadataBuilder) NameMapping()
iceberg.NameMapping {
}
func (b *MetadataBuilder) TrimMetadataLogs(maxEntries int) *MetadataBuilder {
- if len(b.metadataLog) <= maxEntries {
- return b
- }
-
- b.metadataLog = b.metadataLog[len(b.metadataLog)-maxEntries:]
+ b.metadataLog = trimMetadataLog(b.metadataLog, maxEntries)
return b
}
+// trimMetadataLog returns the newest maxEntries entries of log, leaving log
itself
+// unchanged.
+func trimMetadataLog(log []MetadataLogEntry, maxEntries int)
[]MetadataLogEntry {
+ if len(log) <= maxEntries {
+ return log
+ }
+
+ return log[len(log)-maxEntries:]
+}
+
func (b *MetadataBuilder) AppendMetadataLog(entry MetadataLogEntry)
*MetadataBuilder {
b.metadataLog = append(b.metadataLog, entry)
diff --git a/table/metadata_builder_internal_test.go
b/table/metadata_builder_internal_test.go
index cfd6b56f9..78b79298e 100644
--- a/table/metadata_builder_internal_test.go
+++ b/table/metadata_builder_internal_test.go
@@ -1489,6 +1489,60 @@ func TestExpireMetadataLog(t *testing.T) {
require.Len(t, meta.(*metadataV2).MetadataLog, 2)
}
+func TestBuildDoesNotGrowMetadataLog(t *testing.T) {
+ builder := builderWithoutChanges(2)
+ require.NoError(t, builder.SetProperties(map[string]string{"test.prop":
"value"}))
+
+ for range 3 {
+ meta, err := builder.Build()
+ require.NoError(t, err)
+ require.Len(t, meta.(*metadataV2).MetadataLog, 1)
+ require.Empty(t, builder.metadataLog, "Build must not append to
the builder's own log")
+ }
+}
+
+func TestBuildWithoutChangesAddsNoMetadataLogEntry(t *testing.T) {
+ builder := builderWithoutChanges(2)
+ require.NotNil(t, builder.previousFileEntry, "setup: the builder must
have a previous file")
+
+ for range 2 {
+ meta, err := builder.Build()
+ require.NoError(t, err)
+ require.Empty(t, meta.(*metadataV2).MetadataLog)
+ }
+}
+
+func TestBuildDoesNotTrimTheBuilderMetadataLog(t *testing.T) {
+ base := builderWithoutChanges(2)
+ require.NoError(t,
base.SetProperties(map[string]string{MetadataPreviousVersionsMaxKey: "2"}))
+ meta1, err := base.Build()
+ require.NoError(t, err)
+
+ builder2, err := MetadataBuilderFromBase(meta1,
"s3://bucket/test/location/metadata/v2.json")
+ require.NoError(t, err)
+ require.NoError(t,
builder2.SetProperties(map[string]string{"test.prop": "value1"}))
+ meta2, err := builder2.Build()
+ require.NoError(t, err)
+
+ builder3, err := MetadataBuilderFromBase(meta2,
"s3://bucket/test/location/metadata/v3.json")
+ require.NoError(t, err)
+ require.NoError(t,
builder3.SetProperties(map[string]string{"test.prop": "value2"}))
+
+ want := []string{
+ "s3://bucket/test/location/metadata/v2.json",
+ "s3://bucket/test/location/metadata/v3.json",
+ }
+ for range 2 {
+ meta, err := builder3.Build()
+ require.NoError(t, err)
+ var got []string
+ for entry := range meta.PreviousFiles() {
+ got = append(got, entry.MetadataFile)
+ }
+ require.Equal(t, want, got)
+ }
+}
+
func TestMetadataLogTrimsWithUpdatedProperty(t *testing.T) {
// This test validates that when write.metadata.previous-versions-max is
// updated during a metadata build, the trimming logic uses the NEW
value