KranzL opened a new issue, #2040: URL: https://github.com/apache/iceberg-go/issues/2040
### Proposed change `Transaction.RewriteDataFiles` runs compaction groups one at a time. The atomic path loops over groups at table/rewrite_data_files.go:335-361 and calls `ExecuteCompactionGroup` at :344. The partial-progress path has the same shape per batch at :626-647 and calls it at :631. Each group is an independent read, decode, encode and write pipeline. Only intra-group read concurrency exists today, through `WithCompactionScanConcurrency` (table/rewrite_data_files.go:249-252); writes use the clustered single-writer path (:466-468). I propose a `MaxConcurrency` field on `RewriteDataFilesOptions` (table/rewrite_data_files.go:180-220, next to `MaxCommits` and `PartialProgress`) that bounds how many `ExecuteCompactionGroup` calls run at once in both paths. Zero and one keep today's sequential behavior. A negative value is rejected the way `MaxCommits` below zero is at :572-573. The following hold on main at 833e10c: - Output names cannot collide. Each `WriteRecords` call mints a fresh random write UUID when `args.writeUUID` is nil (table/arrow_utils.go:2076-2078 and :2218-2220), so two groups never share an output name. - Scans are independent per call. `Table.Scan` has a value receiver and builds a fresh `Scan` per call (table/table.go:1338). The per-call scans share the snapshot manifest cache from #1970 (:1344), which is built for concurrent scans. - The mechanism has an in-repo precedent. Orphan cleanup bounds parallel work with errgroup `SetLimit` (table/orphan_cleanup.go:562-563) with a `runtime.GOMAXPROCS` default (:284) set through `WithCleanupMaxConcurrency` (:132-138). No open issue tracks this. `gh issue list -R apache/iceberg-go --state open --search "compaction concurrent parallel groups"` returns nothing. The broader "compaction parallel" search returns only #860 (an Overwrite race) and #1178 (the REST scan planning epic), neither about compaction groups. ### Measured evidence Setup: a partitioned v2 table with 8 partitions and 8 data files per partition, 30000 rows per file of (int64 id, string data, 96 hex char payload, float64 score), 222.8 MB on disk. Each group holds one partition (8 files, 240000 rows). Group count N runs `RewriteDataFiles` over the first N groups of the same table with default options (the atomic path) and a fresh transaction per N; the table is never committed between runs, so every N sees identical inputs. Timing covers the `RewriteDataFiles` call only, including the per-group `ExecuteCompactionGroup` calls and the in-transaction rewrite commit. User CPU is the RUSAGE_SELF utime delta around the call. Apple M3 Pro, 11 cores, go1.25.9 darwin/arm64, `go test -run TestZZGroupScaleSequential -count=5`, wall in ms, five runs: | groups | wall, five runs | user CPU in s, five runs | user CPU per wall | |---:|---|---|---| | 1 | 311, 232, 214, 218, 225 | 0.274, 0.272, 0.255, 0.266, 0.274 | 0.88, 1.17, 1.19, 1.22, 1.22 | | 2 | 455, 436, 439, 428, 440 | 0.529, 0.524, 0.514, 0.504, 0.534 | 1.16, 1.20, 1.17, 1.18, 1.21 | | 4 | 826, 856, 812, 839, 845 | 0.989, 1.029, 0.998, 1.005, 1.015 | 1.20, 1.20, 1.23, 1.20, 1.20 | | 8 | 1667, 1633, 1695, 1693, 1677 | 1.961, 1.982, 2.026, 2.036, 2.031 | 1.18, 1.21, 1.20, 1.20, 1.21 | Wall time is linear in group count: 8 groups cost 7.5x one group (medians 225 ms and 1677 ms) while user CPU stays near one core on an 11 core machine. Per-group Parquet encode is single threaded, so the other cores stay idle as groups are added. ### Design `MaxConcurrency` is an int on `RewriteDataFilesOptions`, not a `CompactionGroupOption`. `GroupOptions` are forwarded to every `ExecuteCompactionGroup` call (table/rewrite_data_files.go:215-219) and tune one group's pipeline; the number of groups in flight is a property of the executor loop. `MaxCommits` (:186-190) is the precedent for an executor-level int on this struct. Both executor loops run `ExecuteCompactionGroup` calls under an errgroup with `SetLimit(MaxConcurrency)`, following table/orphan_cleanup.go:562-563. Results are collected per index and applied in original group order before the existing staging and commit logic, so commit semantics do not change. Memory sizing: peak record-pipeline memory is `MaxConcurrency` times the per-group bound PR #2039 states, which is `workers x (rows in the largest task + n) + (recordBatchBufferSize + 2) x n` rows, where `workers` is the scan worker count and `n` is the `WithCompactionArrowBatchSize` value. PR #2039 is open; its docs are at table/rewrite_data_files.go:249-283 on branch `bounded-scan-reorder-heap`. The PR that implements this issue repeats that formula times `MaxConcurrency` in the `MaxConcurrency` doc. ### Correctness requirements - Group results are applied in original group order (`rewrite.ApplyResult` and the batch commit lists follow the input order), so manifests and `RewriteResult` counters are deterministic. - The first group error cancels the remaining groups and is returned. In the atomic path no group result is applied after a failure. - Partial-progress cleanup removes outputs from every group in the batch, including a group that failed mid-write. `ExecuteCompactionGroup` returns partial `NewDataFiles` with its error (:500-507), and the batch cleanup (:612-624) must see those files from every group, failed or not. - Commit semantics are unchanged: one snapshot for the atomic path, at most `MaxCommits` batch snapshots for partial progress. - The default is unchanged: zero and one mean sequential execution with today's behavior. ### Validation plan - Unit tests for ordered application of concurrent group results, first-error cancellation, and partial-progress cleanup covering groups that failed mid-write. - `go test -race` over the rewrite suites and the rolling, fanout and clustered writer suites, the set the #1993 review ran. - A multi-group benchmark at concurrency 1, 2, 4 and 8 on a table of the shape in the measured evidence section, showing wall scaling with group count at each level. ### Related - #1993 added `WithCompactionArrowBatchSize`, `WithCompactionRecordBatchBufferSize` and `WithCompactionParquetRowGroupLimit`. This issue scales the pipeline those knobs bound. - #2038 and PR #2039 bound the per-group scan reorder heap. The `MaxConcurrency` doc multiplies the bound #2039 states. - #1970 shares decoded manifest lists across concurrent scans of one table. Concurrent groups rely on that cache. ### Willingness to contribute - [x] I can contribute this improvement/feature independently - [ ] I would be willing to contribute this improvement/feature with guidance from the Iceberg community - [ ] I cannot contribute this improvement/feature at this time ### Specifications - [x] Table - [ ] View - [ ] REST - [ ] Puffin - [ ] Encryption - [ ] Other -- 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]
