-
Notifications
You must be signed in to change notification settings - Fork 246
fix(table): skip empty historical-delete manifest merges #2057
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -507,18 +507,32 @@ 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 | ||
| // 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 !writerClosed { | ||
| if wr != nil && !writerClosed { | ||
| internal.CheckedClose(wr, &err) | ||
| } | ||
| 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) { | ||
| if err != nil { | ||
|
|
@@ -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,13 @@ 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 { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Java's |
||
| 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 +607,9 @@ func (m *manifestMergeManager) mergeGroup(firstManifest iceberg.ManifestFile, sp | |
| if err != nil { | ||
| return nil, err | ||
| } | ||
| output = append(output, created) | ||
| if created != nil { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Separate from this fix but related: the |
||
| output = append(output, created) | ||
| } | ||
| } | ||
|
|
||
| return output, nil | ||
|
|
@@ -667,7 +699,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 +1795,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) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This |
||
|
|
||
| return []Update{addSnap, retainingSnapshotRef}, []Requirement{AssertRefSnapshotID(branch, baseHeadID)}, nil | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -578,6 +578,130 @@ 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 range 10 { | ||
| 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() | ||
| 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) { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The three new tests cover the happy paths well, but #2037's validation list explicitly calls for "a read error occurring before an output writer is created," and none of these hit it. I'd add a two-manifest bin where the first is all-historical-delete (writer stays |
||
| 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 +2022,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() | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
This assigns into the named-return
err, but every call site checks the closure's return value inside afor entry, err := range ...where the:=shadows the outererr, so the explicitreturn nil, erris what actually propagates and this write to the outererris dead. It reads like the named return is tracked automatically, which invites a future edit to drop the explicit check on the false assumptionerris already set. I'd give the closure its own local: