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]

Reply via email to