Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 33 additions & 0 deletions data_file_codec_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,37 @@ func TestMarshalAvroEntryDecimalPartitionRoundTrip(t *testing.T) {
require.True(t, got.Equals(DecimalLiteral(want)))
}

func TestMarshalAvroEntryTimestampPartitionRoundTrip(t *testing.T) {
// Micro cases run first.
// A nano type sharing their schema-cache key would decode as Timestamp instead of TimestampNano.
for _, tc := range []struct {
typ Type
want any
}{
{TimestampType{}, Timestamp(1_700_000_000_000_000)},
{TimestampTzType{}, Timestamp(1_700_000_000_000_000)},
{TimestampNsType{}, TimestampNano(1_700_000_000_000_000_123)},
{TimestampTzNsType{}, TimestampNano(1_700_000_000_000_000_123)},
} {
t.Run(tc.typ.String(), func(t *testing.T) {
schema := NewSchema(0, NestedField{ID: 1, Name: "ts", Type: tc.typ})
spec := NewPartitionSpecID(1, PartitionField{SourceIDs: []int{1}, FieldID: 1000, Name: "ts", Transform: IdentityTransform{}})

path := "s3://bucket/ns/tbl/data/ts.parquet"
builder, err := NewDataFileBuilder(spec, EntryContentData, path, ParquetFile, map[int]any{1000: tc.want}, nil, nil, 1, 1024)
require.NoError(t, err)
df, ok := builder.Build().(*dataFile)
require.True(t, ok)

encoded, err := df.MarshalAvroEntry(spec, schema, 3)
require.NoError(t, err)
decoded, err := unmarshalAvroDataFileEntry(encoded, spec, schema, 3)
require.NoError(t, err)
require.Equal(t, tc.want, decoded.Partition()[1000])
})
}
}

// snapshotAvroFields returns a deep copy of every avro-tagged field on
// d, keyed by field name. Slices, maps, byte arrays, and pointer
// targets are reconstructed so the snapshot is fully independent of d
Expand Down Expand Up @@ -357,6 +388,8 @@ func TestManifestEntrySchemaForMatchesPartitionAvroShape(t *testing.T) {
TimeType{},
TimestampType{},
TimestampTzType{},
TimestampNsType{},
TimestampTzNsType{},
UUIDType{},
BooleanType{},
BinaryType{},
Expand Down
214 changes: 176 additions & 38 deletions table/rewrite_data_files.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import (
"github.com/apache/iceberg-go"
iceberginternal "github.com/apache/iceberg-go/internal"
iceio "github.com/apache/iceberg-go/io"
"golang.org/x/sync/errgroup"
)

// RewriteResult summarizes a completed compaction.
Expand Down Expand Up @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

}

// CompactionGroupOption configures a single [ExecuteCompactionGroup]
Expand Down Expand Up @@ -252,7 +277,10 @@ func WithCompactionTargetFileSize(size int64) CompactionGroupOption {
// The scan runs min(n, number of tasks) workers, and each worker holds
// the decoded batches of at most one task until the writer has taken
// them, so the worker count multiplies the read-side term of the
// memory bound stated on [WithCompactionArrowBatchSize].
// memory bound stated on [WithCompactionArrowBatchSize]. It also
// multiplies [RewriteDataFilesOptions.MaxConcurrentGroups] for file
// opens: N groups scan with up to n workers each, so about N times n
// files are open at once.
func WithCompactionScanConcurrency(n int) CompactionGroupOption {
return func(c *compactionGroupConfig) {
c.scanConcurrency = n
Expand Down Expand Up @@ -337,6 +365,9 @@ func (t *Transaction) RewriteDataFiles(ctx context.Context, groups []CompactionT
if _, err := t.txnMeta(); err != nil {
return nil, err
}
if opts.MaxConcurrentGroups < 0 {
return nil, fmt.Errorf("%w: MaxConcurrentGroups must be non-negative", ErrInvalidOperation)
}
if opts.PartialProgress {
return t.rewriteDataFilesPartial(ctx, groups, opts)
}
Expand All @@ -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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This doesn't fully close r4108885244. The handle is opened with the caller's ctx, and the blob-backed IOs keep that ctx: blobfs.New(ctx, …) (reached via io.LoadFSFunc) stores it, and Remove calls bfs.Delete(bfs.ctx, key) (io/gocloud/blobfs/blob.go:309-315). Once the caller cancels, every Remove in cleanupAtomicRewriteOutputs fails with context.Canceled, so on S3/GCS/Azure the cancel path still orphans all outputs. I reproduced it with a test IO whose Remove honors the ctx passed to fsF: a sequential run cancelled after group 0 wrote goes from 4 parquet files before to 5 after, and the error ends with clean up atomic rewrite outputs: remove …: context canceled.

Open the cleanup handle detached from cancellation. The partial path has the same pattern at :746 (outside this diff); please change it there too.

Suggested change
fs, err := t.tbl.fsF(ctx)
fs, err := t.tbl.fsF(context.WithoutCancel(ctx))

CleanupAfterCancelUsesOpenFS can't catch this because blockPathIO.Remove is plain LocalFS and ignores ctx. Have cancelOnDoneFSF (rewrite_data_files_test.go:1806) wrap the IO so Remove fails once the ctx it was created with is done; then the test pins the real behavior.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

if err != nil {
return result, fmt.Errorf("open table IO for atomic rewrite: %w", err)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

}

gr, err := ExecuteCompactionGroup(ctx, t.tbl, group, opts.GroupOptions...)
var applied []CompactionGroupResult
if opts.MaxConcurrentGroups > 1 {
results, err := executeCompactionGroups(ctx, t.tbl, groups, opts.GroupOptions, opts.MaxConcurrentGroups)
if err != nil {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The concurrent atomic path leaks output files on failure, and the concurrency makes it strictly worse than the sequential loop.

When a group fails here we return without touching results, so every group that already finished writing its compacted parquet is orphaned on disk (nothing gets committed to a manifest). Sequentially that's bounded to the groups before the failing one; with MaxConcurrency, up to N groups race ahead and finish real writes before cancellation propagates, so the leak scales with the knob.

The partial-progress branch already handles this via cleanupBatch(err, results...). I'd mirror it here: on error call cleanupCompactionOutputs(fs, results) (fs from t.tbl.fsF(ctx)) before returning, and do the same in the sequential else branch. wdyt?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in 30ffe36. The concurrent branch passes its results to cleanupAtomicRewriteOutputs at table/rewrite_data_files.go:382, and the sequential branch tracks applied results and cleans them at :391 and :400. The helper at :453 opens the IO through t.tbl.fsF and joins a cleanup failure with the cause instead of replacing it.

return result, err
return result, cleanupAtomicRewriteOutputs(fs, results, err)
}

if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 {
continue
applied = results
for _, gr := range results {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This per-group apply block (skip empty, ApplyResult, accumulateGroupMetrics, stage SafePosDeletes/SafeDeletionVectors) is now copy-pasted four times: atomic concurrent and sequential here, plus the two partial-progress branches around :700. A future tweak to the accounting has to land in all four or the concurrent and sequential paths silently diverge, which is the exact thing MaxConcurrency's doc promises won't happen.

Could we factor the per-group apply into a small helper shared by both branches?

func applyGroupResult(rewrite *Rewrite, result *RewriteResult, staged map[string]struct{}, gr CompactionGroupResult) { ... }

Same shape for the partial accumulation. Thoughts?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in 30ffe36. Both atomic branches call applyAtomicGroupResult at table/rewrite_data_files.go:439, and both partial branches call appendPartialGroupResult at :843.

applyAtomicGroupResult(rewrite, result, stagedDeleteFiles, gr)
}
} else {
for _, group := range groups {
if err := ctx.Err(); err != nil {
return result, cleanupAtomicRewriteOutputs(fs, applied, err)
}

rewrite.ApplyResult(gr)
accumulateGroupMetrics(result, gr)
for _, df := range gr.SafePosDeletes {
stagedDeleteFiles[df.FilePath()] = struct{}{}
}
for _, df := range gr.SafeDeletionVectors {
stagedDeleteFiles[df.FilePath()] = struct{}{}
if len(group.Tasks) == 0 {
continue
}

gr, err := ExecuteCompactionGroup(ctx, t.tbl, group, opts.GroupOptions...)
if err != nil {
return result, cleanupAtomicRewriteOutputs(fs, append(applied, gr), err)
}
applied = append(applied, gr)

applyAtomicGroupResult(rewrite, result, stagedDeleteFiles, gr)
}
}

Expand Down Expand Up @@ -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))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

}

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
}
Comment on lines +474 to +477

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

main enabled the modernize linter in #2072 after this branch's last CI run. Merged with current main, modernize reports rewrite_data_files.go:475:5: if statement can be modernized using max, so the golangci-lint step fails once you rebase. #2072 rewrote the same pattern in table/internal/variant_shredding.go.

Suggested change
limit := min(maxConcurrentGroups, len(groups))
if limit < 1 {
limit = 1
}
limit := max(min(maxConcurrentGroups, len(groups)), 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 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

if len(group.Tasks) == 0 {
continue
}
g.Go(func() error {
gr, err := ExecuteCompactionGroup(runCtx, tbl, group, groupOpts...)
results[i] = gr
groupErrs[i] = err
if err != nil {
slog.Warn("compaction group failed", "index", i, "err", err)
cancelRuns()
Comment on lines +492 to +494

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: every group that gets cancelled because a sibling failed also logs here, so one real failure with N groups in flight produces N compaction group failed warnings, N-1 of them context canceled. Skip the log (still calling cancelRuns()) when errors.Is(err, context.Canceled) && runCtx.Err() != nil, or log those at Debug.

}

return err
})
}
if err := g.Wait(); err != nil {
var firstErr, firstNonContextErr error
for _, groupErr := range groupErrs {
if groupErr == nil {
continue
}
if firstErr == nil {
firstErr = groupErr
}
if !errors.Is(groupErr, context.Canceled) && !errors.Is(groupErr, context.DeadlineExceeded) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

runCtx is a plain WithCancel, so cancelling the remaining groups only ever produces context.Canceled. A DeadlineExceeded here is either the caller's deadline, which the substitution at :522 already handles, or a real timeout inside that group (SDK and http.Client timeouts wrap it). Treating it as sibling cancellation hides the real failure: with group 2 failing on open …: object store read timeout: context deadline exceeded and the caller ctx still live, this returns write compacted files for group "p0": context canceled, while a sequential run returns group 2's timeout. That's what's left of r4108885251, and it contradicts the "failures match a sequential run" doc at :226-230.

Suggested change
if !errors.Is(groupErr, context.Canceled) && !errors.Is(groupErr, context.DeadlineExceeded) {
if !errors.Is(groupErr, context.Canceled) {

firstNonContextErr = groupErr

break
}
}
selected := firstNonContextErr
if selected == nil {
selected = firstErr
}
if selected == nil {
selected = err
}
if ctxErr := ctx.Err(); ctxErr != nil && (errors.Is(selected, context.Canceled) || errors.Is(selected, context.DeadlineExceeded)) {
return results, ctxErr
}

return results, selected
}

return results, nil
}

// ExecuteCompactionGroup reads a compaction group's tasks (with
// deletes applied), writes consolidated output files via
// [WriteRecords], and computes the position-delete files safe to
Expand Down Expand Up @@ -639,26 +761,26 @@ func (t *Transaction) rewriteDataFilesPartial(ctx context.Context, groups []Comp
return cause
}

for _, group := range batchGroups {
if err := ctx.Err(); err != nil {
return result, cleanupBatch(err)
}

gr, err := ExecuteCompactionGroup(ctx, current, group, opts.GroupOptions...)
if opts.MaxConcurrentGroups > 1 {
results, err := executeCompactionGroups(ctx, current, batchGroups, opts.GroupOptions, opts.MaxConcurrentGroups)
Comment on lines +764 to +765

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

In partial-progress mode this only runs groups concurrently inside one commit batch, and batches still run one after another with a catalog commit in between. A batch is ceil(groups/MaxCommits) groups (:731), so with the default MaxCommits (10) and 10 or fewer groups every batch holds one group and MaxConcurrentGroups does nothing. With 8 groups, PartialProgress, MaxConcurrentGroups: 4 and scan concurrency 1, peak concurrent data-file opens is 1 at MaxCommits: 0 and 3-4 at MaxCommits: 1, which is why the partial cases in the tests set MaxCommits: 1. Java's max-concurrent-file-group-rewrites doesn't depend on commit batching. At minimum, state this cap in the MaxConcurrentGroups doc (:222-244); overlapping execution across batches can be a follow-up.

if err != nil {
return result, cleanupBatch(err, gr)
return result, cleanupBatch(err, results...)
}

if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 {
continue
for _, gr := range results {
batchResults, rewrittenFiles = appendPartialGroupResult(batchResults, rewrittenPaths, rewrittenFiles, gr)
}
batchResults = append(batchResults, gr)
for _, df := range gr.OldDataFiles {
if _, ok := rewrittenPaths[df.FilePath()]; ok {
continue
} else {
for _, group := range batchGroups {
if err := ctx.Err(); err != nil {
return result, cleanupBatch(err)
}
rewrittenPaths[df.FilePath()] = struct{}{}
rewrittenFiles = append(rewrittenFiles, df)

gr, err := ExecuteCompactionGroup(ctx, current, group, opts.GroupOptions...)
if err != nil {
return result, cleanupBatch(err, gr)
}

batchResults, rewrittenFiles = appendPartialGroupResult(batchResults, rewrittenPaths, rewrittenFiles, gr)
}
}

Expand Down Expand Up @@ -752,6 +874,22 @@ func (t *Transaction) rewriteDataFilesPartial(ctx context.Context, groups []Comp
return result, nil
}

func appendPartialGroupResult(batchResults []CompactionGroupResult, rewrittenPaths map[string]struct{}, rewrittenFiles []iceberg.DataFile, gr CompactionGroupResult) ([]CompactionGroupResult, []iceberg.DataFile) {
if len(gr.OldDataFiles) == 0 && len(gr.NewDataFiles) == 0 {
return batchResults, rewrittenFiles
}
batchResults = append(batchResults, gr)
for _, df := range gr.OldDataFiles {
if _, ok := rewrittenPaths[df.FilePath()]; ok {
continue
}
rewrittenPaths[df.FilePath()] = struct{}{}
rewrittenFiles = append(rewrittenFiles, df)
}

return batchResults, rewrittenFiles
}

func recordCommittedRewriteBatch(
result *RewriteResult,
next *Table,
Expand Down
Loading
Loading