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]

Reply via email to