From c3975900a2585f53f0bfba385616604927ca3af6 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 24 Sep 2026 23:11:26 -0700 Subject: [PATCH 1/3] fix(table): skip empty historical-delete manifest merges Signed-off-by: 1fanwang <1fannnw@gmail.com> --- table/snapshot_producers.go | 67 +++++++++++++++++------- table/snapshot_producers_test.go | 90 +++++++++++++++++++++++++++++++- 2 files changed, 136 insertions(+), 21 deletions(-) diff --git a/table/snapshot_producers.go b/table/snapshot_producers.go index 490fb2528..07bb5f99a 100644 --- a/table/snapshot_producers.go +++ b/table/snapshot_producers.go @@ -507,17 +507,31 @@ func (m *manifestMergeManager) createManifest(specID int, bin []iceberg.Manifest return nil, err } - wr, path, counter, fileCloser, err := m.snap.newManifestWriter(spec) - if err != nil { - return nil, err - } - defer internal.CheckedClose(fileCloser, &err) + var wr *iceberg.ManifestWriter + var path string + var counter *internal.CountingWriter + var fileCloser io.Closer writerClosed := false defer func() { - if !writerClosed { + if wr != nil && !writerClosed { internal.CheckedClose(wr, &err) } }() + defer func() { + if fileCloser != nil { + internal.CheckedClose(fileCloser, &err) + } + }() + + ensureWriter := func() error { + if wr != nil { + return nil + } + + wr, path, counter, fileCloser, err = m.snap.newManifestWriter(spec) + + return err + } for _, manifest := range bin { for entry, err := range m.snap.iterManifestEntries(manifest, false) { @@ -527,12 +541,21 @@ func (m *manifestMergeManager) createManifest(specID int, bin []iceberg.Manifest switch { case entry.Status() == iceberg.EntryStatusDELETED && entry.SnapshotID() == m.snap.snapshotID: + if err = ensureWriter(); err != nil { + return nil, err + } // only files deleted by this snapshot should be added to the new manifest err = wr.Delete(entry) case entry.Status() == iceberg.EntryStatusADDED && entry.SnapshotID() == m.snap.snapshotID: + if err = ensureWriter(); err != nil { + return nil, err + } // added entries from this snapshot are still added, otherwise they should be existing err = wr.Add(entry) case entry.Status() != iceberg.EntryStatusDELETED: + if err = ensureWriter(); err != nil { + return nil, err + } // add all non-deleted files from the old manifest as existing files err = wr.Existing(entry) } @@ -543,6 +566,10 @@ func (m *manifestMergeManager) createManifest(specID int, bin []iceberg.Manifest } } + if wr == nil { + return nil, nil + } + // close the writer to force a flush and ensure counter.Count is accurate writerClosed = true if err := wr.Close(); err != nil { @@ -577,7 +604,9 @@ func (m *manifestMergeManager) mergeGroup(firstManifest iceberg.ManifestFile, sp if err != nil { return nil, err } - output = append(output, created) + if created != nil { + output = append(output, created) + } } return output, nil @@ -667,7 +696,8 @@ func newMergeAppendFilesProducer(op Operation, txn *Transaction, fs iceio.WriteF minCountToMerge: txn.meta.props.GetInt(ManifestMinMergeCountKey, ManifestMinMergeCountDefault), mergeEnabled: txn.meta.props.GetBool(ManifestMergeEnabledKey, ManifestMergeEnabledDefault), mergeConcurrency: manifestMergeConcurrencyLimit( - txn.meta.props.GetInt(ManifestMergeMaxConcurrencyKey, ManifestMergeMaxConcurrencyDefault)), + txn.meta.props.GetInt(ManifestMergeMaxConcurrencyKey, ManifestMergeMaxConcurrencyDefault), + ), } return prod @@ -1762,16 +1792,13 @@ func (sp *snapshotProducer) commitManifests(newManifests, addedContent []iceberg // creates it). baseHeadID := sp.txn.baseRefSnapshotID(branch) - return []Update{ - addSnap, - // Carry over the branch's existing retention settings so advancing - // the ref on commit does not silently discard them. The update - // encodes exactly the current ref's retention (settings the branch - // lacks stay 0 and are dropped by the `omitempty` tags); the catalog - // applies a set-snapshot-ref as a pure replace, so this fully - // determines the resulting ref rather than merging with the old one. - sp.txn.meta.NewRetainingSnapshotRefUpdate(branch, sp.snapshotID, BranchRef), - }, []Requirement{ - AssertRefSnapshotID(branch, baseHeadID), - }, nil + // Carry over the branch's existing retention settings so advancing + // the ref on commit does not silently discard them. The update + // encodes exactly the current ref's retention (settings the branch + // lacks stay 0 and are dropped by the `omitempty` tags); the catalog + // applies a set-snapshot-ref as a pure replace, so this fully + // determines the resulting ref rather than merging with the old one. + retainingSnapshotRef := sp.txn.meta.NewRetainingSnapshotRefUpdate(branch, sp.snapshotID, BranchRef) + + return []Update{addSnap, retainingSnapshotRef}, []Requirement{AssertRefSnapshotID(branch, baseHeadID)}, nil } diff --git a/table/snapshot_producers_test.go b/table/snapshot_producers_test.go index 47e60e13d..1370f4347 100644 --- a/table/snapshot_producers_test.go +++ b/table/snapshot_producers_test.go @@ -578,6 +578,93 @@ func TestManifestMergeManagerClosesWriterOnError(t *testing.T) { require.ErrorIs(t, err, errLimitedWrite) } +func TestManifestMergeSkipsHistoricalDeletedOnlyManifest(t *testing.T) { + spec := iceberg.NewPartitionSpec() + schema := simpleSchema() + trackIO := newTrackingIO() + txn := createTestTransaction(t, trackIO, spec) + sp := newFastAppendFilesProducer(OpAppend, txn, trackIO, nil, nil) + + oldSnapshotID := sp.snapshotID - 1 + sequenceNumber := int64(1) + df := newTestDataFile(t, spec, "file://deleted.parquet", nil) + manifestFile := writeTestManifestWithContent(t, trackIO, spec, schema, oldSnapshotID, + "table-location/metadata/historical-delete.avro", iceberg.ManifestContentData, + []iceberg.ManifestEntry{ + iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, &oldSnapshotID, &sequenceNumber, nil, df), + }) + + trackIO.writers = make(map[string]*trackingWriteCloser) + + mgr := manifestMergeManager{snap: sp} + created, err := mgr.createManifest(spec.ID(), []iceberg.ManifestFile{manifestFile}) + require.NoError(t, err) + require.Nil(t, created) + require.Zero(t, trackIO.GetWriterCount()) +} + +func TestManifestMergeKeepsCurrentSnapshotDeletedEntries(t *testing.T) { + spec := iceberg.NewPartitionSpec() + schema := simpleSchema() + txn, wfs := createTestTransactionWithMemIO(t, spec) + sp := newFastAppendFilesProducer(OpAppend, txn, wfs, nil, nil) + + sequenceNumber := int64(1) + df := newTestDataFile(t, spec, "file://deleted.parquet", nil) + manifestFile := writeTestManifestWithContent(t, wfs, spec, schema, sp.snapshotID, + "mem://default/table-location/metadata/current-delete.avro", iceberg.ManifestContentData, + []iceberg.ManifestEntry{ + iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, &sp.snapshotID, &sequenceNumber, nil, df), + }) + + mgr := manifestMergeManager{snap: sp} + created, err := mgr.createManifest(spec.ID(), []iceberg.ManifestFile{manifestFile}) + require.NoError(t, err) + require.NotNil(t, created) + + var entries []iceberg.ManifestEntry + for entry, err := range created.Entries(wfs, false) { + require.NoError(t, err) + entries = append(entries, entry) + } + require.Len(t, entries, 1) + require.Equal(t, iceberg.EntryStatusDELETED, entries[0].Status()) + require.Equal(t, sp.snapshotID, entries[0].SnapshotID()) + require.Equal(t, df.FilePath(), entries[0].DataFile().FilePath()) +} + +func TestManifestMergeGroupDropsEmptyMergedBin(t *testing.T) { + spec := iceberg.NewPartitionSpec() + schema := simpleSchema() + txn, wfs := createTestTransactionWithMemIO(t, spec) + sp := newFastAppendFilesProducer(OpAppend, txn, wfs, nil, nil) + + oldSnapshotID := sp.snapshotID - 1 + sequenceNumber := int64(1) + manifests := make([]iceberg.ManifestFile, 0, 2) + for i := range 2 { + df := newTestDataFile(t, spec, "file://deleted-"+strconv.Itoa(i)+".parquet", nil) + manifests = append(manifests, writeTestManifestWithContent(t, wfs, spec, schema, oldSnapshotID, + "mem://default/table-location/metadata/historical-delete-"+strconv.Itoa(i)+".avro", + iceberg.ManifestContentData, + []iceberg.ManifestEntry{ + iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, &oldSnapshotID, &sequenceNumber, nil, df), + })) + } + + mgr := manifestMergeManager{ + targetSizeBytes: manifests[0].Length() + manifests[1].Length(), + minCountToMerge: 1, + mergeEnabled: true, + mergeConcurrency: 1, + snap: sp, + } + + merged, err := mgr.mergeGroup(manifests[0], spec.ID(), manifests) + require.NoError(t, err) + require.Empty(t, merged) +} + func TestOverwriteFilesExistingManifestsClosesWriterOnError(t *testing.T) { spec := partitionedSpec() schema := simpleSchema() @@ -1898,7 +1985,8 @@ func newTestDeletionVectorForRef(t *testing.T, spec iceberg.PartitionSpec, path, builder, err := iceberg.NewDataFileBuilder( spec, iceberg.EntryContentPosDeletes, path, iceberg.PuffinFile, - nil, nil, nil, 1, 1) + nil, nil, nil, 1, 1, + ) require.NoError(t, err, "new deletion vector builder") return builder.ReferencedDataFile(referencedDataFile).Build() From a0fc20558159f253139eefdd1b677d6d310b384b Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Mon, 28 Sep 2026 14:24:44 -0700 Subject: [PATCH 2/3] fix(table): close ManifestWriter before file on error path Collapse the two deferred closes into a single closure with explicit order so a mid-loop write failure flushes into an open file, not a closed one. Add a regression test that fails a later entry write after the writer is opened. Document the intentional zero-count-manifest divergence from Java/PyIceberg. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- table/snapshot_producers.go | 7 ++++-- table/snapshot_producers_test.go | 38 ++++++++++++++++++++++++++++++++ 2 files changed, 43 insertions(+), 2 deletions(-) diff --git a/table/snapshot_producers.go b/table/snapshot_producers.go index 07bb5f99a..777c0b677 100644 --- a/table/snapshot_producers.go +++ b/table/snapshot_producers.go @@ -512,12 +512,12 @@ func (m *manifestMergeManager) createManifest(specID int, bin []iceberg.Manifest var counter *internal.CountingWriter var fileCloser io.Closer writerClosed := false + // Close the ManifestWriter before the underlying file so a mid-loop + // write failure flushes into an open file, not a closed one. defer func() { if wr != nil && !writerClosed { internal.CheckedClose(wr, &err) } - }() - defer func() { if fileCloser != nil { internal.CheckedClose(fileCloser, &err) } @@ -566,6 +566,9 @@ func (m *manifestMergeManager) createManifest(specID int, bin []iceberg.Manifest } } + // A bin with no live entries produces no manifest. This diverges from + // Java/PyIceberg, which write a zero-count manifest; the omission is + // intentional and should not be "fixed" later. if wr == nil { return nil, nil } diff --git a/table/snapshot_producers_test.go b/table/snapshot_producers_test.go index 1370f4347..40508c520 100644 --- a/table/snapshot_producers_test.go +++ b/table/snapshot_producers_test.go @@ -578,6 +578,44 @@ func TestManifestMergeManagerClosesWriterOnError(t *testing.T) { require.ErrorIs(t, err, errLimitedWrite) } +func TestManifestMergeManagerClosesWriterBeforeFileOnWriteFailure(t *testing.T) { + spec := iceberg.NewPartitionSpec() + schema := simpleSchema() + + // Use a byte-limited IO that fails after the writer is opened and + // has written some data, but before all entries are processed. + mem := newMemIO(manifestHeaderSize(t, 2, spec, schema), errLimitedWrite) + txn := createTestTransaction(t, mem, spec) + + sp := newFastAppendFilesProducer(OpAppend, txn, mem, nil, nil) + df := newTestDataFile(t, spec, "file://data-1.parquet", nil) + + // Build a manifest with enough entries that the writer is opened + // and multiple writes occur before the failure. + entries := make([]iceberg.ManifestEntry, 0, 10) + for i := 0; i < 10; i++ { + entries = append(entries, iceberg.NewManifestEntry( + iceberg.EntryStatusADDED, &sp.snapshotID, nil, nil, df)) + } + + manifestPath := "table-location/metadata/manifest-1.avro" + var manifestBuf bytes.Buffer + manifestFile, err := iceberg.WriteManifest(manifestPath, &manifestBuf, 2, spec, schema, sp.snapshotID, entries) + require.NoError(t, err, "write manifest") + require.NoError(t, mem.WriteFile(manifestPath, manifestBuf.Bytes())) + + mgr := manifestMergeManager{snap: sp} + _, err = mgr.createManifest(spec.ID(), []iceberg.ManifestFile{manifestFile}) + require.ErrorIs(t, err, errLimitedWrite) + + // The writer must be closed before the file closer. If it were not, + // the ManifestWriter.Close() flush would write to an already-closed + // file and return a "write after close" error. The fact that we + // get errLimitedWrite (not a write-after-close error) confirms + // the ordering is correct. +} + + func TestManifestMergeSkipsHistoricalDeletedOnlyManifest(t *testing.T) { spec := iceberg.NewPartitionSpec() schema := simpleSchema() From 84232e4077d723547a7f864615890a05f6ac283c Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Mon, 28 Sep 2026 14:59:52 -0700 Subject: [PATCH 3/3] fix(table): gofmt and intrange lint fixes Signed-off-by: 1fanwang <1fannnw@gmail.com> --- table/snapshot_producers_test.go | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/table/snapshot_producers_test.go b/table/snapshot_producers_test.go index 40508c520..d63806a37 100644 --- a/table/snapshot_producers_test.go +++ b/table/snapshot_producers_test.go @@ -593,7 +593,7 @@ func TestManifestMergeManagerClosesWriterBeforeFileOnWriteFailure(t *testing.T) // Build a manifest with enough entries that the writer is opened // and multiple writes occur before the failure. entries := make([]iceberg.ManifestEntry, 0, 10) - for i := 0; i < 10; i++ { + for range 10 { entries = append(entries, iceberg.NewManifestEntry( iceberg.EntryStatusADDED, &sp.snapshotID, nil, nil, df)) } @@ -615,7 +615,6 @@ func TestManifestMergeManagerClosesWriterBeforeFileOnWriteFailure(t *testing.T) // the ordering is correct. } - func TestManifestMergeSkipsHistoricalDeletedOnlyManifest(t *testing.T) { spec := iceberg.NewPartitionSpec() schema := simpleSchema()