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]