zeroshade commented on code in PR #2041:
URL: https://github.com/apache/iceberg-go/pull/2041#discussion_r4158463257
##########
table/rewrite_data_files.go:
##########
@@ -348,31 +379,38 @@ func (t *Transaction) RewriteDataFiles(ctx
context.Context, groups []CompactionT
rewrite := t.NewRewrite(opts.SnapshotProps)
stagedDeleteFiles := make(map[string]struct{})
- for _, group := range groups {
- if err := ctx.Err(); err != nil {
- return result, err
- }
-
- if len(group.Tasks) == 0 {
- continue
- }
+ fs, err := t.tbl.fsF(ctx)
Review Comment:
This doesn't fully close
[r4108885244](https://github.com/apache/iceberg-go/pull/2041#discussion_r4108885244).
The handle is opened with the caller's `ctx`, and the blob-backed IOs keep
that ctx: `blobfs.New(ctx, …)` (reached via `io.LoadFSFunc`) stores it, and
`Remove` calls `bfs.Delete(bfs.ctx, key)` (io/gocloud/blobfs/blob.go:309-315).
Once the caller cancels, every `Remove` in `cleanupAtomicRewriteOutputs` fails
with `context.Canceled`, so on S3/GCS/Azure the cancel path still orphans all
outputs. I reproduced it with a test IO whose `Remove` honors the ctx passed to
`fsF`: a sequential run cancelled after group 0 wrote goes from 4 parquet files
before to 5 after, and the error ends with `clean up atomic rewrite outputs:
remove …: context canceled`.
Open the cleanup handle detached from cancellation. The partial path has the
same pattern at :746 (outside this diff); please change it there too.
```suggestion
fs, err := t.tbl.fsF(context.WithoutCancel(ctx))
```
`CleanupAfterCancelUsesOpenFS` can't catch this because `blockPathIO.Remove`
is plain `LocalFS` and ignores ctx. Have `cancelOnDoneFSF`
(rewrite_data_files_test.go:1806) wrap the IO so `Remove` fails once the ctx it
was created with is done; then the test pins the real behavior.
##########
table/rewrite_data_files.go:
##########
@@ -401,12 +439,96 @@ func (t *Transaction) RewriteDataFiles(ctx
context.Context, groups []CompactionT
}
if err := rewrite.Commit(ctx); err != nil {
- return result, fmt.Errorf("commit compaction: %w", err)
+ return result, cleanupAtomicRewriteOutputs(fs, applied,
fmt.Errorf("commit compaction: %w", err))
}
return result, nil
}
+func applyAtomicGroupResult(rewrite *RewriteFiles, result *RewriteResult,
stagedDeleteFiles map[string]struct{}, gr CompactionGroupResult) {
+ if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 {
+ return
+ }
+ rewrite.ApplyResult(gr)
+ accumulateGroupMetrics(result, gr)
+ for _, df := range gr.SafePosDeletes {
+ stagedDeleteFiles[df.FilePath()] = struct{}{}
+ }
+ for _, df := range gr.SafeDeletionVectors {
+ stagedDeleteFiles[df.FilePath()] = struct{}{}
+ }
+}
+
+func cleanupAtomicRewriteOutputs(fs iceio.IO, results []CompactionGroupResult,
cause error) error {
+ if err := cleanupCompactionOutputs(fs, results); err != nil {
+ return errors.Join(cause, fmt.Errorf("clean up atomic rewrite
outputs: %w", err))
+ }
+
+ return cause
+}
+
+func executeCompactionGroups(ctx context.Context, tbl *Table, groups
[]CompactionTaskGroup, groupOpts []CompactionGroupOption, maxConcurrentGroups
int) ([]CompactionGroupResult, error) {
+ if err := ctx.Err(); err != nil {
+ return nil, err
+ }
+ limit := min(maxConcurrentGroups, len(groups))
+ if limit < 1 {
+ limit = 1
+ }
+ var g errgroup.Group
+ g.SetLimit(limit)
+ runCtx, cancelRuns := context.WithCancel(ctx)
+ defer cancelRuns()
+ results := make([]CompactionGroupResult, len(groups))
+ groupErrs := make([]error, len(groups))
+ for i, group := range groups {
+ if len(group.Tasks) == 0 {
+ continue
+ }
+ g.Go(func() error {
+ gr, err := ExecuteCompactionGroup(runCtx, tbl, group,
groupOpts...)
+ results[i] = gr
+ groupErrs[i] = err
+ if err != nil {
+ slog.Warn("compaction group failed", "index",
i, "err", err)
+ cancelRuns()
+ }
+
+ return err
+ })
+ }
+ if err := g.Wait(); err != nil {
+ var firstErr, firstNonContextErr error
+ for _, groupErr := range groupErrs {
+ if groupErr == nil {
+ continue
+ }
+ if firstErr == nil {
+ firstErr = groupErr
+ }
+ if !errors.Is(groupErr, context.Canceled) &&
!errors.Is(groupErr, context.DeadlineExceeded) {
Review Comment:
`runCtx` is a plain `WithCancel`, so cancelling the remaining groups only
ever produces `context.Canceled`. A `DeadlineExceeded` here is either the
caller's deadline, which the substitution at :522 already handles, or a real
timeout inside that group (SDK and `http.Client` timeouts wrap it). Treating it
as sibling cancellation hides the real failure: with group 2 failing on `open
…: object store read timeout: context deadline exceeded` and the caller ctx
still live, this returns `write compacted files for group "p0": context
canceled`, while a sequential run returns group 2's timeout. That's what's left
of
[r4108885251](https://github.com/apache/iceberg-go/pull/2041#discussion_r4108885251),
and it contradicts the "failures match a sequential run" doc at :226-230.
```suggestion
if !errors.Is(groupErr, context.Canceled) {
```
##########
table/rewrite_data_files.go:
##########
@@ -401,12 +439,96 @@ func (t *Transaction) RewriteDataFiles(ctx
context.Context, groups []CompactionT
}
if err := rewrite.Commit(ctx); err != nil {
- return result, fmt.Errorf("commit compaction: %w", err)
+ return result, cleanupAtomicRewriteOutputs(fs, applied,
fmt.Errorf("commit compaction: %w", err))
}
return result, nil
}
+func applyAtomicGroupResult(rewrite *RewriteFiles, result *RewriteResult,
stagedDeleteFiles map[string]struct{}, gr CompactionGroupResult) {
+ if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 {
+ return
+ }
+ rewrite.ApplyResult(gr)
+ accumulateGroupMetrics(result, gr)
+ for _, df := range gr.SafePosDeletes {
+ stagedDeleteFiles[df.FilePath()] = struct{}{}
+ }
+ for _, df := range gr.SafeDeletionVectors {
+ stagedDeleteFiles[df.FilePath()] = struct{}{}
+ }
+}
+
+func cleanupAtomicRewriteOutputs(fs iceio.IO, results []CompactionGroupResult,
cause error) error {
+ if err := cleanupCompactionOutputs(fs, results); err != nil {
+ return errors.Join(cause, fmt.Errorf("clean up atomic rewrite
outputs: %w", err))
+ }
+
+ return cause
+}
+
+func executeCompactionGroups(ctx context.Context, tbl *Table, groups
[]CompactionTaskGroup, groupOpts []CompactionGroupOption, maxConcurrentGroups
int) ([]CompactionGroupResult, error) {
+ if err := ctx.Err(); err != nil {
+ return nil, err
+ }
+ limit := min(maxConcurrentGroups, len(groups))
+ if limit < 1 {
+ limit = 1
+ }
Review Comment:
main enabled the `modernize` linter in #2072 after this branch's last CI
run. Merged with current main, modernize reports `rewrite_data_files.go:475:5:
if statement can be modernized using max`, so the golangci-lint step fails once
you rebase. #2072 rewrote the same pattern in
table/internal/variant_shredding.go.
```suggestion
limit := max(min(maxConcurrentGroups, len(groups)), 1)
```
##########
table/rewrite_data_files.go:
##########
@@ -639,26 +761,26 @@ func (t *Transaction) rewriteDataFilesPartial(ctx
context.Context, groups []Comp
return cause
}
- for _, group := range batchGroups {
- if err := ctx.Err(); err != nil {
- return result, cleanupBatch(err)
- }
-
- gr, err := ExecuteCompactionGroup(ctx, current, group,
opts.GroupOptions...)
+ if opts.MaxConcurrentGroups > 1 {
+ results, err := executeCompactionGroups(ctx, current,
batchGroups, opts.GroupOptions, opts.MaxConcurrentGroups)
Review Comment:
In partial-progress mode this only runs groups concurrently inside one
commit batch, and batches still run one after another with a catalog commit in
between. A batch is `ceil(groups/MaxCommits)` groups (:731), so with the
default `MaxCommits` (10) and 10 or fewer groups every batch holds one group
and `MaxConcurrentGroups` does nothing. With 8 groups, `PartialProgress`,
`MaxConcurrentGroups: 4` and scan concurrency 1, peak concurrent data-file
opens is 1 at `MaxCommits: 0` and 3-4 at `MaxCommits: 1`, which is why the
partial cases in the tests set `MaxCommits: 1`. Java's
`max-concurrent-file-group-rewrites` doesn't depend on commit batching. At
minimum, state this cap in the `MaxConcurrentGroups` doc (:222-244);
overlapping execution across batches can be a follow-up.
##########
table/rewrite_data_files.go:
##########
@@ -401,12 +439,96 @@ func (t *Transaction) RewriteDataFiles(ctx
context.Context, groups []CompactionT
}
if err := rewrite.Commit(ctx); err != nil {
- return result, fmt.Errorf("commit compaction: %w", err)
+ return result, cleanupAtomicRewriteOutputs(fs, applied,
fmt.Errorf("commit compaction: %w", err))
}
return result, nil
}
+func applyAtomicGroupResult(rewrite *RewriteFiles, result *RewriteResult,
stagedDeleteFiles map[string]struct{}, gr CompactionGroupResult) {
+ if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 {
+ return
+ }
+ rewrite.ApplyResult(gr)
+ accumulateGroupMetrics(result, gr)
+ for _, df := range gr.SafePosDeletes {
+ stagedDeleteFiles[df.FilePath()] = struct{}{}
+ }
+ for _, df := range gr.SafeDeletionVectors {
+ stagedDeleteFiles[df.FilePath()] = struct{}{}
+ }
+}
+
+func cleanupAtomicRewriteOutputs(fs iceio.IO, results []CompactionGroupResult,
cause error) error {
+ if err := cleanupCompactionOutputs(fs, results); err != nil {
+ return errors.Join(cause, fmt.Errorf("clean up atomic rewrite
outputs: %w", err))
+ }
+
+ return cause
+}
+
+func executeCompactionGroups(ctx context.Context, tbl *Table, groups
[]CompactionTaskGroup, groupOpts []CompactionGroupOption, maxConcurrentGroups
int) ([]CompactionGroupResult, error) {
+ if err := ctx.Err(); err != nil {
+ return nil, err
+ }
+ limit := min(maxConcurrentGroups, len(groups))
+ if limit < 1 {
+ limit = 1
+ }
+ var g errgroup.Group
+ g.SetLimit(limit)
+ runCtx, cancelRuns := context.WithCancel(ctx)
+ defer cancelRuns()
+ results := make([]CompactionGroupResult, len(groups))
+ groupErrs := make([]error, len(groups))
+ for i, group := range groups {
+ if len(group.Tasks) == 0 {
+ continue
+ }
+ g.Go(func() error {
+ gr, err := ExecuteCompactionGroup(runCtx, tbl, group,
groupOpts...)
+ results[i] = gr
+ groupErrs[i] = err
+ if err != nil {
+ slog.Warn("compaction group failed", "index",
i, "err", err)
+ cancelRuns()
Review Comment:
Nit: every group that gets cancelled because a sibling failed also logs
here, so one real failure with N groups in flight produces N `compaction group
failed` warnings, N-1 of them `context canceled`. Skip the log (still calling
`cancelRuns()`) when `errors.Is(err, context.Canceled) && runCtx.Err() != nil`,
or log those at Debug.
##########
table/rewrite_data_files_test.go:
##########
@@ -1401,3 +1426,749 @@ func appendEqualityDelete(t *testing.T, tbl
*table.Table, equalityFieldIDs []int
return out
}
+
+func newMaxConcPartitionedTable(t *testing.T, fs iceio.IO) *table.Table {
+ t.Helper()
+
+ return newMaxConcPartitionedTableWithFSF(t, func(context.Context)
(iceio.IO, error) { return fs, nil })
+}
+
+func newMaxConcPartitionedTableWithFSF(t *testing.T, fsF func(context.Context)
(iceio.IO, error)) *table.Table {
+ t.Helper()
+
+ location := filepath.ToSlash(t.TempDir())
+ schema := iceberg.NewSchema(0,
+ iceberg.NestedField{ID: 1, Name: "id", Type:
iceberg.PrimitiveTypes.Int64, Required: true},
+ iceberg.NestedField{ID: 2, Name: "data", Type:
iceberg.PrimitiveTypes.String, Required: false},
+ )
+ spec := iceberg.NewPartitionSpec(iceberg.PartitionField{
+ SourceIDs: []int{2}, FieldID: 1000, Transform:
iceberg.IdentityTransform{}, Name: "data",
+ })
+ meta, err := table.NewMetadata(schema, &spec, table.UnsortedSortOrder,
location,
+ iceberg.Properties{table.PropertyFormatVersion: "2"})
+ require.NoError(t, err)
+
+ cat := &partialProgressCatalog{metadata: meta}
+
+ return table.New(
+ table.Identifier{"db", "max_conc_test"},
+ meta, location+"/metadata/v1.metadata.json",
+ fsF,
+ cat,
+ )
+}
+
+func addMaxConcPartitions(t *testing.T, tbl *table.Table, partitions,
filesPerPartition, rowsPerFile int) *table.Table {
+ t.Helper()
+
+ var nextID int64 = 1
+ for p := range partitions {
+ partition := fmt.Sprintf("p%d", p)
+ for f := range filesPerPartition {
+ ids := make([]int64, rowsPerFile)
+ for r := range rowsPerFile {
+ ids[r] = nextID
+ nextID++
+ }
+ tbl = addPartitionedRowsOnRef(t, tbl, table.MainBranch,
fmt.Sprintf("p%d-%d", p, f), partition, ids...)
+ }
+ }
+
+ return tbl
+}
+
+func groupsByPartition(t *testing.T, tbl *table.Table)
[]table.CompactionTaskGroup {
+ t.Helper()
+
+ tasks, err := tbl.Scan().PlanFiles(t.Context())
+ require.NoError(t, err)
+
+ byPart := make(map[string][]table.FileScanTask)
+ for _, task := range tasks {
+ part, ok := task.File.Partition()[1000].(string)
+ require.True(t, ok)
+ byPart[part] = append(byPart[part], task)
+ }
+ keys := make([]string, 0, len(byPart))
+ for k := range byPart {
+ keys = append(keys, k)
+ }
+ slices.Sort(keys)
+
+ groups := make([]table.CompactionTaskGroup, 0, len(keys))
+ for _, k := range keys {
+ var total int64
+ for _, task := range byPart[k] {
+ total += task.File.FileSizeBytes()
+ }
+ groups = append(groups, table.CompactionTaskGroup{
+ PartitionKey: k,
+ Tasks: byPart[k],
+ TotalSizeBytes: total,
+ })
+ }
+
+ return groups
+}
+
+func idsByPartitionValue(t *testing.T, tbl *table.Table) map[string][]int64 {
+ t.Helper()
+
+ _, itr, err := tbl.Scan().ToArrowRecords(t.Context())
+ require.NoError(t, err)
+
+ out := make(map[string][]int64)
+ for rec, err := range itr {
+ require.NoError(t, err)
+ dataIdx := rec.Schema().FieldIndices("data")
+ require.NotEmpty(t, dataIdx)
+ idIdx := rec.Schema().FieldIndices("id")
+ require.NotEmpty(t, idIdx)
+ dataCol, ok := rec.Column(dataIdx[0]).(*array.String)
+ require.True(t, ok)
+ idCol, ok := rec.Column(idIdx[0]).(*array.Int64)
+ require.True(t, ok)
+ for i := range int(rec.NumRows()) {
+ out[dataCol.Value(i)] = append(out[dataCol.Value(i)],
idCol.Value(i))
+ }
+ rec.Release()
+ }
+
+ return out
+}
+
+func manifestDataPartitions(t *testing.T, tbl *table.Table) []string {
+ t.Helper()
+
+ snap := tbl.CurrentSnapshot()
+ require.NotNil(t, snap)
+ fs, err := tbl.FS(t.Context())
+ require.NoError(t, err)
+ manifests, err := snap.Manifests(fs)
+ require.NoError(t, err)
+
+ var parts []string
+ for _, m := range manifests {
+ for e, err := range m.Entries(fs, false) {
+ require.NoError(t, err)
+ if e.Status() == iceberg.EntryStatusDELETED {
+ continue
+ }
+ df := e.DataFile()
+ if df.ContentType() != iceberg.EntryContentData {
+ continue
+ }
+ part, ok := df.Partition()[1000].(string)
+ require.True(t, ok)
+ parts = append(parts, part)
+ }
+ }
+
+ return parts
+}
+
+func manifestLiveDataPaths(t *testing.T, tbl *table.Table) []string {
+ t.Helper()
+
+ snap := tbl.CurrentSnapshot()
+ require.NotNil(t, snap)
+ fs, err := tbl.FS(t.Context())
+ require.NoError(t, err)
+ manifests, err := snap.Manifests(fs)
+ require.NoError(t, err)
+
+ var paths []string
+ for _, m := range manifests {
+ for e, err := range m.Entries(fs, false) {
+ require.NoError(t, err)
+ if e.Status() == iceberg.EntryStatusDELETED {
+ continue
+ }
+ df := e.DataFile()
+ if df.ContentType() != iceberg.EntryContentData {
+ continue
+ }
+ paths = append(paths, df.FilePath())
+ }
+ }
+
+ return paths
+}
+
+func TestRewriteDataFiles_MaxConcurrentGroupsMatchesSequential(t *testing.T) {
+ tblSeq := newMaxConcPartitionedTable(t, iceio.LocalFS{})
+ tblSeq = addMaxConcPartitions(t, tblSeq, 8, 2, 5)
+ tblConc := newMaxConcPartitionedTable(t, iceio.LocalFS{})
+ tblConc = addMaxConcPartitions(t, tblConc, 8, 2, 5)
+
+ groupsSeq := groupsByPartition(t, tblSeq)
+ groupsConc := groupsByPartition(t, tblConc)
+ require.Len(t, groupsSeq, 8)
+ require.Len(t, groupsConc, 8)
+
+ txSeq := tblSeq.NewTransaction()
+ resSeq, err := txSeq.RewriteDataFiles(t.Context(), groupsSeq,
table.RewriteDataFilesOptions{})
+ require.NoError(t, err)
+ committedSeq, err := txSeq.Commit(t.Context())
+ require.NoError(t, err)
+
+ txConc := tblConc.NewTransaction()
+ resConc, err := txConc.RewriteDataFiles(t.Context(), groupsConc,
table.RewriteDataFilesOptions{MaxConcurrentGroups: 4})
+ require.NoError(t, err)
+ committedConc, err := txConc.Commit(t.Context())
+ require.NoError(t, err)
+
+ assert.Equal(t, resSeq.RewrittenGroups, resConc.RewrittenGroups)
+ assert.Equal(t, resSeq.AddedDataFiles, resConc.AddedDataFiles)
+ assert.Equal(t, resSeq.RemovedDataFiles, resConc.RemovedDataFiles)
+ assert.Equal(t, resSeq.RemovedPositionDeleteFiles,
resConc.RemovedPositionDeleteFiles)
+ assert.Equal(t, resSeq.RemovedEqualityDeleteFiles,
resConc.RemovedEqualityDeleteFiles)
+ assert.Equal(t, resSeq.RemovedDeletionVectorFiles,
resConc.RemovedDeletionVectorFiles)
+ assert.Equal(t, resSeq.BytesBefore, resConc.BytesBefore)
+ assert.Equal(t, 8, resConc.RewrittenGroups)
+ assert.Equal(t, 16, resConc.RemovedDataFiles)
+ assert.Equal(t, 8, resConc.AddedDataFiles)
+
+ idsSeq := idsByPartitionValue(t, committedSeq)
+ idsConc := idsByPartitionValue(t, committedConc)
+ require.Len(t, idsConc, 8)
+ for p := range 8 {
+ key := fmt.Sprintf("p%d", p)
+ assert.ElementsMatch(t, idsSeq[key], idsConc[key])
+ assert.Len(t, idsConc[key], 10)
+ }
+
+ paths := manifestLiveDataPaths(t, committedConc)
+ require.Len(t, paths, 8)
+ assert.Len(t, map[string]struct{}{paths[0]: {}, paths[1]: {}, paths[2]:
{}, paths[3]: {}, paths[4]: {}, paths[5]: {}, paths[6]: {}, paths[7]: {}}, 8)
+ onDisk := allParquetFiles(t, committedConc.Location())
+ for _, p := range paths {
+ assert.Contains(t, onDisk, p)
+ }
+}
+
+func TestRewriteDataFiles_MaxConcurrentGroupsNegativeRejected(t *testing.T) {
+ tbl := newRewriteTestTable(t)
+
+ tx := tbl.NewTransaction()
+ _, err := tx.RewriteDataFiles(t.Context(), nil,
table.RewriteDataFilesOptions{MaxConcurrentGroups: -1})
+ require.ErrorIs(t, err, table.ErrInvalidOperation)
+
+ txPartial := tbl.NewTransaction()
+ _, err = txPartial.RewriteDataFiles(t.Context(), nil,
table.RewriteDataFilesOptions{PartialProgress: true, MaxConcurrentGroups: -1})
+ require.ErrorIs(t, err, table.ErrInvalidOperation)
+}
+
+func TestRewriteDataFiles_MaxConcurrentGroupsDeterministicOrder(t *testing.T) {
+ tblA := newMaxConcPartitionedTable(t, iceio.LocalFS{})
+ tblA = addMaxConcPartitions(t, tblA, 8, 1, 5)
+ tblB := newMaxConcPartitionedTable(t, iceio.LocalFS{})
+ tblB = addMaxConcPartitions(t, tblB, 8, 1, 5)
+
+ groupsA := groupsByPartition(t, tblA)
+ groupsB := groupsByPartition(t, tblB)
+
+ txA := tblA.NewTransaction()
+ _, err := txA.RewriteDataFiles(t.Context(), groupsA,
table.RewriteDataFilesOptions{MaxConcurrentGroups: 4})
+ require.NoError(t, err)
+ committedA, err := txA.Commit(t.Context())
+ require.NoError(t, err)
+
+ txB := tblB.NewTransaction()
+ _, err = txB.RewriteDataFiles(t.Context(), groupsB,
table.RewriteDataFilesOptions{MaxConcurrentGroups: 4})
+ require.NoError(t, err)
+ committedB, err := txB.Commit(t.Context())
+ require.NoError(t, err)
+
+ orderA := manifestDataPartitions(t, committedA)
+ orderB := manifestDataPartitions(t, committedB)
+ require.Len(t, orderA, 8)
+ require.Len(t, orderB, 8)
+ assert.Equal(t, orderA, orderB)
+ assert.Equal(t, []string{"p0", "p1", "p2", "p3", "p4", "p5", "p6",
"p7"}, orderA)
+}
+
+type failOpenIO struct {
Review Comment:
Nit, non-blocking: `failOpenIO`, `blockPathIO`, `orderFailIO` and
`gateOpenIO` all intercept `Open` by path substring; one IO with an `onOpen
func(name string) error` hook would replace the four.
`newMaxConcPartitionedTableWithFSF` duplicates
`newPartialProgressPartitionedTable` (:871) except for the identifier and fsF,
`allParquetFiles` can replace `parquetFiles`, `manifestDataPartitions` and
`manifestLiveDataPaths` are the same loop, and the bench reimplements
`groupsByPartition`. Folding these together should cut 150+ lines without
losing coverage.
--
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]