laskoviymishka commented on code in PR #2057:
URL: https://github.com/apache/iceberg-go/pull/2057#discussion_r4102962036
##########
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:
This assigns into the named-return `err`, but every call site checks the
closure's return value inside a `for entry, err := range ...` where the `:=`
shadows the outer `err`, so the explicit `return nil, err` is what actually
propagates and this write to the outer `err` is dead. It reads like the named
return is tracked automatically, which invites a future edit to drop the
explicit check on the false assumption `err` is already set. I'd give the
closure its own local:
```go
ensureWriter := func() error {
if wr != nil {
return nil
}
var werr error
wr, path, counter, fileCloser, werr = m.snap.newManifestWriter(spec)
return werr
}
```
##########
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:
Separate from this fix but related: the `len(bin)==1` short-circuit in
`mergeGroup` passes a lone manifest through unfiltered, so a manifest that's
100% historical-DELETED and lands alone in its bin (already at target size, or
gated by `minCountToMerge`) never reaches `createManifest` and carries its dead
tombstones forward indefinitely. This PR fixes the commit crash, not that; fine
as a follow-up, but worth noting so #2037 isn't assumed fully closed.
##########
table/snapshot_producers.go:
##########
@@ -543,6 +566,10 @@ func (m *manifestMergeManager) createManifest(specID int,
bin []iceberg.Manifest
}
}
+ if wr == nil {
Review Comment:
Java's `ManifestMergeManager` and PyIceberg both keep a zero-count manifest
here rather than dropping it, so for identical merge history Go emits fewer
manifest files than the other clients. That's the intended tradeoff from #2037
and it's read-safe, but it's a permanent divergence worth a short comment right
here (and a line in the PR description) so nobody diffing manifest counts
across clients treats it as a bug and "fixes" it back to parity later.
##########
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:
The two cleanup defers run in the wrong order on the error path: LIFO now
closes `fileCloser` before `wr`, so if an entry write fails mid-loop after the
writer's already open, `wr.Close()` flushes the `ManifestWriter` into an
already-closed file. That's the "write after close" hazard
`TestCommitManifestsCloseFailureReturnsNoUpdates` was written to guard, and
every other writer-cleanup site in this file registers the fileCloser defer
first so `wr` flushes first.
I'd collapse both into one closure with an explicit order rather than
relying on LIFO across two statements:
```go
defer func() {
if wr != nil && !writerClosed {
internal.CheckedClose(wr, &err)
}
if fileCloser != nil {
internal.CheckedClose(fileCloser, &err)
}
}()
```
And add a regression test that opens the writer then fails a later entry
write (the `trackingIO`/`failWriteAt` idiom), asserting `NotContains "write
after close"`. None of the new tests hit that ordering today, which is why this
slips through.
##########
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:
This `commitManifests` refactor (and the `NewDataFileBuilder` /
`mergeConcurrency` reformats) are behavior-neutral and unrelated to the
empty-manifest fix. I'd pull them out so blame and bisect stay clean on the
actual change.
##########
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:
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 `nil`) and the second fails to read, and
assert `createManifest` returns the read error, not `(nil, nil)`, and that no
writer was opened. That's exactly the `ensureWriter`-gating seam most likely to
regress if the error check is ever reordered.
--
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]