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]