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


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

Review Comment:
   Wiring cleanup onto `rewrite.Commit` closes part of last round's gap, but 
`rewrite.Commit` only stages the rewrite onto the transaction. The commit that 
actually reaches the catalog and can lose a conflict is `txn.Commit` in 
`Table.RewriteDataFiles` (around line 107), and that path returns the error 
without removing any outputs. Since `Table.RewriteDataFiles` is the turnkey 
entry point and the partial branch already cleans up on `ErrCommitFailed`, the 
atomic action leaks exactly where a real commit conflict lands while partial 
doesn't. I'd clean up in `Table.RewriteDataFiles` when `txn.Commit` fails with 
a provably-uncommitted error (`ErrCommitFailed`), mirroring the partial branch.



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

Review Comment:
   After `cancelRuns()` fires, this loop keeps dispatching every remaining 
group; each runs against the canceled runCtx, fails fast, and logs a Warn, so a 
large plan buries the real error under a wall of context.Canceled warnings. I'd 
`break` out of the launch loop once `runCtx.Err() != nil`.



##########
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:
   Opening fs once up front closes the canceled-open symptom for local 
backends, but not the leak on cloud. `iceio.IO.Remove` takes no ctx, so blobfs 
deletes through the ctx the FileIO was built with (`bfs.Delete(bfs.ctx, key)` 
in `io/gocloud/blobfs/blob.go`), and `fsF(ctx)` threads this request ctx into 
that field. On the cancel/deadline path, which is the main cleanup trigger, 
that stored ctx is already dead, so every cleanup `Remove` fails with 
context.Canceled and the staged outputs stay in object storage. The failure 
just moved from "open FS" to "delete object."
   
   I'd open the cleanup FS with `context.WithoutCancel(ctx)` and use it only 
for the Remove calls, or route cleanup through 
`DeleteFiles(context.WithoutCancel(ctx), paths)` (blobfs threads the per-call 
ctx there) with a bounded timeout. Same shape in the partial path's 
`cleanupBatch`. One caveat: the regression test can't catch this today, since 
`cancelOnDoneFSF` hands back a LocalFS whose Remove ignores ctx, so it passes 
even if cleanup were a no-op on cloud.



##########
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)
+       if err != nil {
+               return result, fmt.Errorf("open table IO for atomic rewrite: 
%w", err)

Review Comment:
   This open is now unconditional, so a rewrite with no non-empty groups, which 
used to return `(result, nil)` without touching IO, now fails here if `fsF` 
errors. I'd skip the open when there's nothing to compact (gate on `len(groups) 
== 0`, or open lazily).



##########
table/rewrite_data_files.go:
##########
@@ -217,6 +218,30 @@ type RewriteDataFilesOptions struct {
        // size, scan concurrency). See the With* helpers returning
        // [CompactionGroupOption].
        GroupOptions []CompactionGroupOption
+
+       // MaxConcurrentGroups bounds how many compaction groups run at once.
+       // Zero and one both mean sequential execution, which is the default.
+       // Larger values run [ExecuteCompactionGroup] calls under a bounded
+       // errgroup and apply their results in the original group order, so
+       // manifests and [RewriteResult] are identical to a sequential run.
+       // A failure cancels the groups still running; the returned error
+       // is the lowest-index failure not caused by that cancellation, or
+       // the lowest-index failure when every group was canceled, so
+       // failures match a sequential run. Every group error is logged.
+       // Peak record-pipeline memory is MaxConcurrentGroups times the
+       // per-group bound stated on [WithCompactionArrowBatchSize]:
+       //
+       //      MaxConcurrentGroups x (workers x (rows in the largest task + n) 
+ (recordBatchBufferSize + 2) x n)
+       //
+       // rows, where workers, n and recordBatchBufferSize are the per-group
+       // values. Multiply rows by the average row width in bytes for a byte
+       // estimate. Delete-side memory is outside this bound. File-open
+       // fan-out multiplies too: every group scans with up to
+       // [WithCompactionScanConcurrency] workers, so N groups open about N
+       // times the scan worker count in files at once. Size the two knobs
+       // together against connection and file descriptor limits. Negative
+       // values are rejected with [ErrInvalidOperation].
+       MaxConcurrentGroups int

Review Comment:
   Worth documenting here that under PartialProgress the groups are split into 
`ceil(len(groups)/MaxCommits)` batches and MaxConcurrentGroups only fans out 
within a batch. With the default MaxCommits=10, any run of 10 or fewer groups 
gets one group per batch and this knob does nothing. I'd call out that 
effective concurrency is `min(MaxConcurrentGroups, ceil(groups/MaxCommits))` so 
partial-mode callers aren't surprised when the speedup disappears.



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