KranzL opened a new pull request, #2041: URL: https://github.com/apache/iceberg-go/pull/2041
This PR depends on #2039 and carries its commits until that PR merges. It adds concurrent compaction groups on top of the bounded scan reorder heap. What changed RewriteDataFilesOptions has a MaxConcurrency field (table/rewrite_data_files.go:222-237). Zero and one keep the sequential loop. Values above one run ExecuteCompactionGroup calls under a bounded errgroup and apply results in the original group order, so manifests and RewriteResult counters match a sequential run. The atomic path branches at :372 and the partial-progress path branches per batch at :712; both call the shared helper at :451-478, which sets the errgroup limit to min(MaxConcurrency, len(groups)). A negative value returns ErrInvalidOperation at :358-359, the same way MaxCommits below zero does. The sizing bound is in the MaxConcurrency doc comment; website/src/api.md:501-515 shows RewriteDataFiles with empty options and lists no options, so that page is unchanged. Tests are TestRewriteDataFiles_MaxConcurrency* in table/rewrite_data_files_test.go:1566-1882: equal committed state and counters vs sequential, deterministic manifest order across runs, group failure returned with partial-progress outputs removed, context cancellation returned, in-flight groups capped at MaxConcurrency and exactly one for 0/1, negative rejected. The benchmark is BenchmarkRewriteDataFilesGroupConcurrency in table/rewrite_data_files_bench_test.go:184-243. Setup builds one partitioned table with 8 partitions of 8 files and 30000 rows each (1,920,000 rows, columns id int64, data string, payload 96-char hex string, score float64, about 223 MB). Each iteration rewrites all 8 groups on a fresh transaction and deletes its outputs afterwards, outside the timer. Measured evidence Command: go test ./table/ -run xxx -bench BenchmarkRewriteDataFilesGroupConcurrency -benchmem -count=5 Hardware: Apple M3 Pro, 11 cores, go1.25.9 darwin/arm64 | MaxConcurrency | ms/op, 5 runs | rows/s median | B/op median | | 1 | 1540, 1518, 1483, 1469, 1470 | 1.29M | 3.17 GB | | 2 | 778, 771, 772, 774, 773 | 2.49M | 3.15 GB | | 4 | 486, 460, 475, 468, 466 | 4.10M | 3.08 GB | | 8 | 399, 389, 383, 388, 385 | 4.95M | 3.06 GB | Median wall time falls from 1483 ms at 1 to 388 ms at 8, a 3.8x speedup. Branch MaxConcurrency=1 sits at 1483 ms against 1459 ms for the sequential-only run of the same fixture on the item 2 branch (2ff58a4, runs 1459, 1463, 1448, 1455, 1461), so the sequential path is unchanged within run noise. B/op is flat near 3.1 GB because every concurrency rewrites the same rows. I also ran a lagging-head check through RewriteDataFiles with 4 groups of 32 files, 4 scan workers per group, and each group's head task gated (scratch test, not committed). MaxConcurrency=1 held 3 of 3 expected tasks at 2.12 MB peak heap delta; MaxConcurrency=4 held 12 of 12 at 8.41 MB, which is 3.97 times the single-group value. The scan-level benchmark on this branch reports 3 gated batches and 2.12, 2.17, 2.38 MB for M 8, 32, 128, so the per-group bound from #2039 holds per group and multiplies by MaxConcurrency. Memory sizing Peak record-pipeline memory is about MaxConcurrency x (workers x (rows in the largest task + n) + (recordBatchBufferSize + 2) x n) rows, times the average row width in bytes. workers, n, and recordBatchBufferSize are the per-group GroupOptions values. Delete-side memory is outside this bound. Correctness The first group error cancels the groups still running through the errgroup context and is returned. The atomic path does not remove outputs on error, which matches its current behavior. The partial-progress path waits for every goroutine in the batch then passes each result, including a failed result carrying partial NewDataFiles, to cleanupBatch before the error returns. Validation go test ./table/ -race -count=1 -run 'RewriteDataFiles|ExecuteCompactionGroup|Rolling|Fanout|Clustered' passes in 9.8s. go test ./table/... -race -count=1 passes (table 30.9s, compaction 11.5s, dv 2.6s, internal 3.6s, substrait 4.1s). make lint reports 0 issues; go vet and gofmt are clean. Fixes #2040. -- 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]
