laskoviymishka commented on code in PR #1913:
URL: https://github.com/apache/iceberg-go/pull/1913#discussion_r3873464310
##########
table/dv_scan_planning_test.go:
##########
@@ -141,6 +141,29 @@ func TestManifestEntries_DVClassification(t *testing.T) {
assert.Len(t, entries.equalityDeleteEntries, 1)
}
+func TestManifestEntries_MergeRejectsUnknownContent(t *testing.T) {
+ snapshotID := int64(1)
+ validEntry := iceberg.NewManifestEntry(iceberg.EntryStatusADDED,
&snapshotID, nil, nil, &mockDataFile{
+ path: "s3://bucket/data/data-001.parquet",
+ contentType: iceberg.EntryContentData,
+ })
+ invalidEntry := iceberg.NewManifestEntry(iceberg.EntryStatusADDED,
&snapshotID, nil, nil, &mockDataFile{
+ path: "s3://bucket/data/unknown.bin",
+ contentType: iceberg.ManifestEntryContent(99),
+ })
+ entryAfterInvalid := iceberg.NewManifestEntry(iceberg.EntryStatusADDED,
&snapshotID, nil, nil, &mockDataFile{
+ path: "s3://bucket/data/data-002.parquet",
+ contentType: iceberg.EntryContentData,
+ })
+
+ entries := newManifestEntries()
+ err := entries.merge([]iceberg.ManifestEntry{validEntry, invalidEntry,
entryAfterInvalid})
Review Comment:
This case only exercises the `kinds == nil` path, since every valid entry
here is data, so classification stays on the fast
`classifiedManifestEntriesForKind` branch and the `kinds != nil` backfill loop
never runs.
I'd add a mixed case like `[data, posDelete, invalid, data]` so the error
hits with `kinds != nil` and the backfill path gets covered, asserting
`dataEntries == [e0]` and `positionalDeleteEntries == [e1]`. While we're here,
worth asserting the other three buckets stay empty too, so a bad entry bleeding
into the wrong bucket can't pass silently.
##########
table/dv_scan_planning_test.go:
##########
@@ -141,6 +141,29 @@ func TestManifestEntries_DVClassification(t *testing.T) {
assert.Len(t, entries.equalityDeleteEntries, 1)
}
+func TestManifestEntries_MergeRejectsUnknownContent(t *testing.T) {
+ snapshotID := int64(1)
+ validEntry := iceberg.NewManifestEntry(iceberg.EntryStatusADDED,
&snapshotID, nil, nil, &mockDataFile{
+ path: "s3://bucket/data/data-001.parquet",
+ contentType: iceberg.EntryContentData,
+ })
+ invalidEntry := iceberg.NewManifestEntry(iceberg.EntryStatusADDED,
&snapshotID, nil, nil, &mockDataFile{
+ path: "s3://bucket/data/unknown.bin",
+ contentType: iceberg.ManifestEntryContent(99),
+ })
+ entryAfterInvalid := iceberg.NewManifestEntry(iceberg.EntryStatusADDED,
&snapshotID, nil, nil, &mockDataFile{
+ path: "s3://bucket/data/data-002.parquet",
+ contentType: iceberg.EntryContentData,
+ })
+
+ entries := newManifestEntries()
+ err := entries.merge([]iceberg.ManifestEntry{validEntry, invalidEntry,
entryAfterInvalid})
+
+ assert.ErrorIs(t, err, ErrInvalidMetadata)
+ assert.ErrorContains(t, err, "unknown DataFileContent type")
+ assert.Equal(t, []iceberg.ManifestEntry{validEntry},
entries.dataEntries)
Review Comment:
One thing worth calling out: this pins the partial-commit-on-error behavior
as expected, but the real caller (`collectManifestEntriesWithSchema`) returns
on the first error and discards these entries, so nothing actually depends on
the partial result landing.
Not asking to change the behavior, just maybe a comment here or in
`classifyManifestEntries` noting the partial commit is incidental, so this test
doesn't get read as a hard requirement to preserve. wdyt?
##########
table/scanner.go:
##########
@@ -147,30 +164,128 @@ func newManifestEntries() *manifestEntries {
}
}
-func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error {
- m.mu.Lock()
- defer m.mu.Unlock()
-
- for _, entry := range entries {
- dataFile := entry.DataFile()
- switch dataFile.ContentType() {
- case iceberg.EntryContentData:
- m.dataEntries = append(m.dataEntries, entry)
- case iceberg.EntryContentPosDeletes:
- if IsDeletionVector(dataFile) {
- m.dvEntries = append(m.dvEntries, entry)
- } else {
- m.positionalDeleteEntries =
append(m.positionalDeleteEntries, entry)
+func classifyManifestEntry(entry iceberg.ManifestEntry) (manifestEntryKind,
error) {
+ dataFile := entry.DataFile()
+ switch dataFile.ContentType() {
+ case iceberg.EntryContentData:
+ return manifestEntryData, nil
+ case iceberg.EntryContentPosDeletes:
+ if IsDeletionVector(dataFile) {
+ return manifestEntryDV, nil
+ }
+
+ return manifestEntryPositionalDelete, nil
+ case iceberg.EntryContentEqDeletes:
+ return manifestEntryEqualityDelete, nil
+ default:
+ return 0, fmt.Errorf("%w: unknown DataFileContent type (%s):
%s",
+ ErrInvalidMetadata, dataFile.ContentType(), entry)
+ }
+}
+
+func newClassifiedManifestEntries(counts [manifestEntryKindCount]int)
classifiedManifestEntries {
+ return classifiedManifestEntries{
+ dataEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryData]),
+ positionalDeleteEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryPositionalDelete]),
+ equalityDeleteEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryEqualityDelete]),
+ dvEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryDV]),
+ }
+}
+
Review Comment:
Minor: the homogeneous path assigns the caller's slice straight into the
result (`classified.dataEntries = entries`) with no copy, so the returned
struct aliases the `openManifest` slice that `merge` was handed. It's safe
today since `merge` only reads it through `append`, but the aliasing is
intentional and undocumented.
A one-line comment would save the next person who tries to mutate one of
these buckets in place from a surprise.
##########
table/scanner.go:
##########
@@ -147,30 +164,128 @@ func newManifestEntries() *manifestEntries {
}
}
-func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error {
- m.mu.Lock()
- defer m.mu.Unlock()
-
- for _, entry := range entries {
- dataFile := entry.DataFile()
- switch dataFile.ContentType() {
- case iceberg.EntryContentData:
- m.dataEntries = append(m.dataEntries, entry)
- case iceberg.EntryContentPosDeletes:
- if IsDeletionVector(dataFile) {
- m.dvEntries = append(m.dvEntries, entry)
- } else {
- m.positionalDeleteEntries =
append(m.positionalDeleteEntries, entry)
+func classifyManifestEntry(entry iceberg.ManifestEntry) (manifestEntryKind,
error) {
+ dataFile := entry.DataFile()
+ switch dataFile.ContentType() {
+ case iceberg.EntryContentData:
+ return manifestEntryData, nil
+ case iceberg.EntryContentPosDeletes:
+ if IsDeletionVector(dataFile) {
+ return manifestEntryDV, nil
+ }
+
+ return manifestEntryPositionalDelete, nil
+ case iceberg.EntryContentEqDeletes:
+ return manifestEntryEqualityDelete, nil
+ default:
+ return 0, fmt.Errorf("%w: unknown DataFileContent type (%s):
%s",
+ ErrInvalidMetadata, dataFile.ContentType(), entry)
+ }
+}
+
+func newClassifiedManifestEntries(counts [manifestEntryKindCount]int)
classifiedManifestEntries {
+ return classifiedManifestEntries{
+ dataEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryData]),
+ positionalDeleteEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryPositionalDelete]),
+ equalityDeleteEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryEqualityDelete]),
+ dvEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryDV]),
+ }
+}
+
+func classifiedManifestEntriesForKind(kind manifestEntryKind, entries
[]iceberg.ManifestEntry) classifiedManifestEntries {
+ classified := classifiedManifestEntries{}
+ switch kind {
+ case manifestEntryData:
+ classified.dataEntries = entries
+ case manifestEntryPositionalDelete:
+ classified.positionalDeleteEntries = entries
+ case manifestEntryEqualityDelete:
+ classified.equalityDeleteEntries = entries
+ case manifestEntryDV:
+ classified.dvEntries = entries
+ }
+
+ return classified
+}
+
+func classifyManifestEntries(entries []iceberg.ManifestEntry)
(classifiedManifestEntries, error) {
+ if len(entries) == 0 {
+ return classifiedManifestEntries{}, nil
+ }
+
+ firstKind, err := classifyManifestEntry(entries[0])
+ if err != nil {
+ return classifiedManifestEntries{}, err
+ }
+
+ var (
+ kinds []manifestEntryKind
+ counts [manifestEntryKindCount]int
+ )
+ for i := 1; i < len(entries); i++ {
+ kind, err := classifyManifestEntry(entries[i])
+ if err != nil {
+ if kinds == nil {
+ return
classifiedManifestEntriesForKind(firstKind, entries[:i]), err
}
- case iceberg.EntryContentEqDeletes:
- m.equalityDeleteEntries =
append(m.equalityDeleteEntries, entry)
- default:
- return fmt.Errorf("%w: unknown DataFileContent type
(%s): %s",
- ErrInvalidMetadata, dataFile.ContentType(),
entry)
+
+ classified := newClassifiedManifestEntries(counts)
+ for j, validEntry := range entries[:i] {
+ appendClassifiedManifestEntry(&classified,
kinds[j], validEntry)
+ }
+
+ return classified, err
+ }
+
+ if kinds == nil && kind != firstKind {
+ kinds = make([]manifestEntryKind, len(entries))
+ for j := range i {
+ kinds[j] = firstKind
+ }
+ counts[firstKind] = i
+ }
+ if kinds != nil {
+ kinds[i] = kind
+ counts[kind]++
}
}
- return nil
+ if kinds == nil {
+ return classifiedManifestEntriesForKind(firstKind, entries), nil
+ }
+
+ classified := newClassifiedManifestEntries(counts)
+ for i, entry := range entries {
+ appendClassifiedManifestEntry(&classified, kinds[i], entry)
+ }
+
+ return classified, nil
+}
+
Review Comment:
Neither this switch nor the one in `classifiedManifestEntriesForKind` has a
default, so if a fifth `manifestEntryKind` ever gets added and someone forgets
to wire it here, the entry is silently dropped with no error. `exhaustive`
isn't enabled, so nothing catches that statically, and the non-data fast paths
aren't covered by a test either, so a drop would be invisible.
Not blocking, but I'd add a `default: panic(fmt.Sprintf("unhandled
manifestEntryKind %d", kind))` to both switches so a future miss fails loudly.
wdyt?
##########
table/scanner.go:
##########
@@ -147,30 +164,128 @@ func newManifestEntries() *manifestEntries {
}
}
-func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error {
- m.mu.Lock()
- defer m.mu.Unlock()
-
- for _, entry := range entries {
- dataFile := entry.DataFile()
- switch dataFile.ContentType() {
- case iceberg.EntryContentData:
- m.dataEntries = append(m.dataEntries, entry)
- case iceberg.EntryContentPosDeletes:
- if IsDeletionVector(dataFile) {
- m.dvEntries = append(m.dvEntries, entry)
- } else {
- m.positionalDeleteEntries =
append(m.positionalDeleteEntries, entry)
+func classifyManifestEntry(entry iceberg.ManifestEntry) (manifestEntryKind,
error) {
+ dataFile := entry.DataFile()
+ switch dataFile.ContentType() {
+ case iceberg.EntryContentData:
+ return manifestEntryData, nil
+ case iceberg.EntryContentPosDeletes:
+ if IsDeletionVector(dataFile) {
+ return manifestEntryDV, nil
+ }
+
+ return manifestEntryPositionalDelete, nil
+ case iceberg.EntryContentEqDeletes:
+ return manifestEntryEqualityDelete, nil
+ default:
+ return 0, fmt.Errorf("%w: unknown DataFileContent type (%s):
%s",
+ ErrInvalidMetadata, dataFile.ContentType(), entry)
+ }
+}
+
+func newClassifiedManifestEntries(counts [manifestEntryKindCount]int)
classifiedManifestEntries {
+ return classifiedManifestEntries{
+ dataEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryData]),
+ positionalDeleteEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryPositionalDelete]),
+ equalityDeleteEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryEqualityDelete]),
+ dvEntries: make([]iceberg.ManifestEntry, 0,
counts[manifestEntryDV]),
+ }
+}
+
+func classifiedManifestEntriesForKind(kind manifestEntryKind, entries
[]iceberg.ManifestEntry) classifiedManifestEntries {
+ classified := classifiedManifestEntries{}
+ switch kind {
+ case manifestEntryData:
+ classified.dataEntries = entries
+ case manifestEntryPositionalDelete:
+ classified.positionalDeleteEntries = entries
+ case manifestEntryEqualityDelete:
+ classified.equalityDeleteEntries = entries
+ case manifestEntryDV:
+ classified.dvEntries = entries
+ }
+
+ return classified
+}
+
+func classifyManifestEntries(entries []iceberg.ManifestEntry)
(classifiedManifestEntries, error) {
+ if len(entries) == 0 {
+ return classifiedManifestEntries{}, nil
+ }
+
+ firstKind, err := classifyManifestEntry(entries[0])
+ if err != nil {
+ return classifiedManifestEntries{}, err
+ }
+
+ var (
+ kinds []manifestEntryKind
+ counts [manifestEntryKindCount]int
+ )
+ for i := 1; i < len(entries); i++ {
+ kind, err := classifyManifestEntry(entries[i])
+ if err != nil {
+ if kinds == nil {
+ return
classifiedManifestEntriesForKind(firstKind, entries[:i]), err
}
- case iceberg.EntryContentEqDeletes:
- m.equalityDeleteEntries =
append(m.equalityDeleteEntries, entry)
- default:
- return fmt.Errorf("%w: unknown DataFileContent type
(%s): %s",
- ErrInvalidMetadata, dataFile.ContentType(),
entry)
+
+ classified := newClassifiedManifestEntries(counts)
+ for j, validEntry := range entries[:i] {
+ appendClassifiedManifestEntry(&classified,
kinds[j], validEntry)
+ }
+
+ return classified, err
+ }
+
+ if kinds == nil && kind != firstKind {
+ kinds = make([]manifestEntryKind, len(entries))
+ for j := range i {
+ kinds[j] = firstKind
+ }
+ counts[firstKind] = i
+ }
+ if kinds != nil {
+ kinds[i] = kind
+ counts[kind]++
}
}
- return nil
+ if kinds == nil {
+ return classifiedManifestEntriesForKind(firstKind, entries), nil
+ }
+
+ classified := newClassifiedManifestEntries(counts)
+ for i, entry := range entries {
+ appendClassifiedManifestEntry(&classified, kinds[i], entry)
+ }
+
+ return classified, nil
+}
+
+func appendClassifiedManifestEntry(classified *classifiedManifestEntries, kind
manifestEntryKind, entry iceberg.ManifestEntry) {
+ switch kind {
+ case manifestEntryData:
+ classified.dataEntries = append(classified.dataEntries, entry)
+ case manifestEntryPositionalDelete:
+ classified.positionalDeleteEntries =
append(classified.positionalDeleteEntries, entry)
+ case manifestEntryEqualityDelete:
+ classified.equalityDeleteEntries =
append(classified.equalityDeleteEntries, entry)
+ case manifestEntryDV:
+ classified.dvEntries = append(classified.dvEntries, entry)
+ }
+}
+
+func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error {
+ classified, err := classifyManifestEntries(entries)
Review Comment:
I'd restore `defer m.mu.Unlock()` here rather than unlocking explicitly at
the end.
Realistically these four appends won't panic, so this isn't a live deadlock,
but it's the one spot that dropped the defer idiom every other method on `m.mu`
still uses, and if any append ever panics (OOM, or a future change) the mutex
stays held and every subsequent `merge` deadlocks. It's free to keep the
guarantee: `m.mu.Lock()` then `defer m.mu.Unlock()` right at the top.
--
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]