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]

Reply via email to