KranzL commented on code in PR #2041:
URL: https://github.com/apache/iceberg-go/pull/2041#discussion_r4139890723


##########
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:
   Done in 6c9a36e. The atomic path tracks applied results across both branches 
and the rewrite.Commit failure branch passes them to 
cleanupAtomicRewriteOutputs at table/rewrite_data_files.go:442. Unconditional 
cleanup is safe here since rewrite.Commit only stages onto the caller 
transaction while the catalog commit happens later in the caller 
Transaction.Commit. CommitFailureCleansOutputs in 1116187 fails the staging 
with an unsupported content type and asserts the before/after parquet set 
matches for maxConc 4 and 0.



##########
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:
   Done in 6c9a36e, second option. The atomic path opens fs once up front at 
table/rewrite_data_files.go:382 and cleanupAtomicRewriteOutputs takes the open 
iceio.IO handle instead of (ctx, tbl), mirroring the partial path. 
CleanupAfterCancelUsesOpenFS in 1116187 runs an fsF that fails on a done 
context and asserts full cleanup after mid-run cancellation on both the 
sequential and concurrent paths.



##########
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:
   Done in 6c9a36e, lowest-index option plus logging. executeCompactionGroups 
records per-group errors, logs every non-nil one via slog.Warn, and returns the 
lowest-index failure not caused by sibling cancellation (falling back to the 
lowest-index failure when every group was canceled), so failures match a 
sequential run; the round-1 context substitution now applies to the selected 
error. Cancellation also broadcasts over a plain WithCancel context instead of 
errgroup.WithContext: with the shared cause, a still-running sibling scan 
iterator surfaced the failing group error value through context.Cause instead 
of context.Canceled, which the new ordering test caught before the broadcast 
change. FailureReturnsLowestIndexError in 1116187 orders two injected failures 
so index 3 finishes first and asserts index 1 error is returned while both 
appear in the log.



-- 
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