laskoviymishka commented on code in PR #2041:
URL: https://github.com/apache/iceberg-go/pull/2041#discussion_r4108885236
##########
table/rewrite_data_files.go:
##########
@@ -348,31 +376,32 @@ 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 opts.MaxConcurrentGroups > 1 {
+ results, err := executeCompactionGroups(ctx, t.tbl, groups,
opts.GroupOptions, opts.MaxConcurrentGroups)
+ if err != nil {
+ return result, cleanupAtomicRewriteOutputs(ctx, t.tbl,
results, err)
Review Comment:
We clean up here when a group fails, but the terminal `rewrite.Commit(ctx)`
around line 432 has no equivalent. If every group writes successfully and then
Commit hits a concurrent-modification or catalog conflict (the common
atomic-rewrite failure), none of the new files get removed, and
`MaxConcurrentGroups` makes it worse since N groups wrote in parallel and one
conflict orphans all N groups' output at once. I'd wire
`cleanupAtomicRewriteOutputs` into the commit-failure branch too, the same way
partial-progress cleans up its own commit failures.
##########
table/rewrite_data_files.go:
##########
@@ -407,6 +436,65 @@ func (t *Transaction) RewriteDataFiles(ctx
context.Context, groups []CompactionT
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(ctx context.Context, tbl *Table, results
[]CompactionGroupResult, cause error) error {
+ fs, err := tbl.fsF(ctx)
+ if err != nil {
+ return errors.Join(cause, fmt.Errorf("open table IO to clean up
atomic rewrite outputs: %w", err))
+ }
+ 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
+ }
+ g, gctx := errgroup.WithContext(ctx)
+ g.SetLimit(limit)
+ results := make([]CompactionGroupResult, len(groups))
+ for i, group := range groups {
+ if len(group.Tasks) == 0 {
+ continue
+ }
+ g.Go(func() error {
+ gr, err := ExecuteCompactionGroup(gctx, tbl, group,
groupOpts...)
+ results[i] = gr
+
+ return err
+ })
+ }
+ if err := g.Wait(); err != nil {
+ if ctxErr := ctx.Err(); ctxErr != nil && (errors.Is(err,
context.Canceled) || errors.Is(err, context.DeadlineExceeded)) {
Review Comment:
On concurrent failure the returned error is whichever goroutine finishes
first, so a genuine group error can get masked by a sibling's
`context.Canceled` (the outer ctx wasn't canceled, so the `ctxErr` guard just
above doesn't fire) and the caller loses the real cause. It's also
nondeterministic versus sequential's always-first-in-list, so the doc line
claiming the result is identical to a sequential run overclaims on the error
path. I'd log every non-nil group error via slog so the root cause stays
visible, and either pick the lowest-index failing group's error or scope that
doc claim to the success path.
##########
table/rewrite_data_files.go:
##########
@@ -407,6 +436,65 @@ func (t *Transaction) RewriteDataFiles(ctx
context.Context, groups []CompactionT
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(ctx context.Context, tbl *Table, results
[]CompactionGroupResult, cause error) error {
+ fs, err := tbl.fsF(ctx)
Review Comment:
`cleanupAtomicRewriteOutputs` re-derives the fs from `ctx`, and on the
cancellation path that ctx is already canceled, so for a real cloud FS
(s3/gocloud resolve their config off the context) `fsF(ctx)` can fail before we
reach `Remove` and we return having deleted nothing. That's the same leak
round-1 was closing, just moved down a level. The partial path avoids it by
opening `fs` once before executing any group and passing the handle into
cleanup; I'd mirror that here and have this take the already-open `iceio.IO`
instead of `(ctx, tbl)`. The stubs don't catch it because every `fsF` just
returns `fs, nil` and ignores ctx.
--
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]