1fanwang commented on code in PR #2057:
URL: https://github.com/apache/iceberg-go/pull/2057#discussion_r4162555275
##########
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() {
Review Comment:
Done in
https://github.com/apache/iceberg-go/commit/4e958128577cf5ba785051ae0e3b6b09f481b8ef.
Both cleanups now run in one deferred closure that closes the writer before
the file.
##########
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)
Review Comment:
Done in
https://github.com/apache/iceberg-go/commit/4e958128577cf5ba785051ae0e3b6b09f481b8ef.
The closure now uses its own local error.
##########
table/snapshot_producers.go:
##########
@@ -543,6 +566,10 @@ func (m *manifestMergeManager) createManifest(specID int,
bin []iceberg.Manifest
}
}
+ if wr == nil {
Review Comment:
Done in
https://github.com/apache/iceberg-go/commit/4e958128577cf5ba785051ae0e3b6b09f481b8ef.
The comment and the PR description now note the divergence from Java.
##########
table/snapshot_producers.go:
##########
@@ -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 {
Review Comment:
Agreed that it's separate. I left singleton bins unchanged here because this
fix only covers bins that are actually merged, so #2037 still needs a follow-up
for the lone tombstone manifest.
##########
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) {
Review Comment:
Done in
https://github.com/apache/iceberg-go/commit/4e958128577cf5ba785051ae0e3b6b09f481b8ef.
The new test has a two-manifest bin whose second manifest fails to read, and
it asserts that the read error comes back with no writer opened.
##########
table/snapshot_producers.go:
##########
@@ -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)
Review Comment:
Done in
https://github.com/apache/iceberg-go/commit/dcc96bbbaa851bbfc2e66eaa507c297fe54b39b5.
The commitManifests reflow and the reformats are gone, so the PR diff only
touches snapshot_producers.go and its tests.
##########
table/snapshot_producers_test.go:
##########
@@ -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)
Review Comment:
Done in
https://github.com/apache/iceberg-go/commit/4e958128577cf5ba785051ae0e3b6b09f481b8ef.
The test now uses trackingIO and asserts that the error has no write after
close, so reversing the close order fails it.
##########
table/snapshot_producers.go:
##########
@@ -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)
Review Comment:
Done in
https://github.com/apache/iceberg-go/commit/4e958128577cf5ba785051ae0e3b6b09f481b8ef.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]