KranzL opened a new issue, #2038:
URL: https://github.com/apache/iceberg-go/issues/2038

   ### Proposed change
   
   The doc on `WithCompactionArrowBatchSize` at 
table/rewrite_data_files.go:258-266 says that together with 
`WithCompactionRecordBatchBufferSize` it "bounds the memory held by the record 
pipeline specifically: buffered batches times rows per batch". That holds only 
when the group's scan reads one task at a time. At the default 
`WithCompactionScanConcurrency` (table/rewrite_data_files.go:249-252, zero 
means runtime.GOMAXPROCS), or any value above one, the scan's reorder step 
retains an unbounded number of decoded batches whenever one task is slower than 
the others.
   
   The mechanism, on main at 833e10c:
   
   - `recordBatchesFromTasksAndDeletes` (table/arrow_scanner.go:2230) starts 
`numWorkers = min(concurrency, len(tasks))` goroutines (:2234) that share a 
`records` channel of capacity `numWorkers` (:2236). Tasks are dispatched in 
index order (:2313-2320), and a worker takes the next task as soon as it 
finishes one. Batches are sent to `records` from `recordsFromTask` (:1895) and 
`processRecordsWithPlans` (:1707).
   - `createIteratorWithCleanup` (table/arrow_scanner.go:2125) wraps `records` 
in `MakeSequencedChanWithDiscard` with `bufferSize = numWorkers` (:2130), so 
batches leave the scan in task order, then record order.
   - `MakeSequencedChanWithDiscard` (table/internal/utils.go:92-113) pushes 
every value it reads from `source` into a priority queue (`heap.Push` at :108) 
and pops only values that are next in order (:109-112). Nothing reads 
`pq.Len()` and nothing stops it from reading `source`.
   
   While task 0 is still opening or decoding its file, the other `numWorkers - 
1` workers finish task 1, task 2, and so on, and every batch they produce lands 
in the heap, because nothing can be emitted before task 0's first batch. The 
heap grows until task 0 produces or the scan runs out of tasks. Its size is set 
by the total decoded size of every other task, not by any tunable.
   
   The approval review on #1993 recorded this as its first non-blocking 
follow-up: 
https://github.com/apache/iceberg-go/pull/1993#pullrequestreview-5270700379. No 
open issue tracks it.
   
   I propose to bound the heap with per-worker in-flight credits in 
`recordBatchesFromTasksAndDeletes`, and to rewrite the doc at 
table/rewrite_data_files.go:258-266 so it states the bound that actually holds.
   
   ### Measured evidence
   
   Reproduction: an unpartitioned v2 table with M data files, each a single row 
group of 2048 rows of (int64, 256 byte string), 536.5 KB per file as Arrow. The 
table's IO is `iceio.LocalFS` wrapped so that `Open` of the data file behind 
task 0 blocks until the benchmark releases it. The scan is 
`tbl.Scan(table.WithMaxConcurrency(4))` followed by `ReadTasks` over all M 
planned tasks, consumed from a goroutine. After the gate is hit, the benchmark 
waits until the scanner has opened and closed the other M - 1 files, sleeps 150 
ms, runs `runtime.GC()`, and reads `runtime.MemStats.HeapInuse`. The reported 
number is that reading minus a reading taken after a GC just before 
`ReadTasks`. The gate then opens and the benchmark checks that all M x 2048 
rows arrive. The control rows run the same scan with `WithMaxConcurrency(1)`, 
where the only worker is the gated one.
   
   Apple M3 Pro, 11 cores, go1.25.9 darwin/arm64, `go test -bench -benchtime=1x 
-count=3`, one sample per count:
   
   | files M | scan concurrency | files fully read while task 0 is gated | 
HeapInuse delta in MB, three runs |
   |---:|---:|---:|---|
   | 8 | 4 | 7 | 4.63, 4.50, 4.65 |
   | 32 | 4 | 31 | 18.98, 19.42, 19.60 |
   | 128 | 4 | 127 | 79.80, 79.07, 79.97 |
   | 8 | 1 | 0 | 0.00, 0.00, 0.00 |
   | 32 | 1 | 0 | 0.00, 0.00, 0.00 |
   | 128 | 1 | 0 | 0.00, 0.00, 0.02 |
   
   At concurrency 4 the retained heap is about 1.2 x (M - 1) x 536.5 KB and 
grows linearly with M: going from 32 to 128 files multiplies the delta by 4.1. 
An earlier invocation of the same benchmark gave 4.48 to 4.59, 19.02 to 19.54 
and 79.37 to 79.56 MB for the concurrency 4 rows. With the default 
`read.parquet.batch-size` of 131072 rows each of these files is one batch, so 
the heap holds M - 1 entries; `WithArrowBatchSize` changes the size of each 
entry, not the count of entries the heap will accept.
   
   For a compaction group this means that whenever one input file is slow to 
open or decode (a cold object store read, a large footer, a task with a large 
positional delete set), the memory held between the scan and the writer is 
every other batch in the group, independent of `WithCompactionArrowBatchSize` 
and `WithCompactionRecordBatchBufferSize`.
   
   ### Design
   
   Per-worker in-flight credits. Each of the `numWorkers` goroutines in 
`recordBatchesFromTasksAndDeletes` may have at most K batches that it has 
produced but that `MakeSequencedChanWithDiscard` has not yet emitted at 
table/internal/utils.go:111 or discarded at close (:99-105). A worker takes a 
credit before sending a data batch to `records` and blocks, also on 
`scanCtx.Done()`, when it has none. The credit returns when the batch leaves 
the reorder step. Error records (`enumeratedRecord{Err: err}`) carry no batch 
and take no credit. The credit belongs to the worker, not to the task index, so 
the worker on the lagging task never waits for credits held by other workers 
and produces as soon as its file opens.
   
   Bound: `records` plus the heap hold at most `numWorkers x K` batches. The 
sequenced output channel (`make(chan T, bufferSize)` at 
table/internal/utils.go:95) holds `numWorkers` more. The writer's input queue 
holds `recordBatchBufferSize` (default 64, 
table/rewrite_data_files.go:275-278). With `WithCompactionArrowBatchSize(n)` 
the record pipeline of one group then holds at most `numWorkers x K + 
numWorkers + recordBatchBufferSize` batches of n rows, plus the batch the 
writer is encoding. That is the sentence the doc at 
table/rewrite_data_files.go:258-266 should carry, and 
`WithCompactionScanConcurrency`'s doc should say that it is the `numWorkers` 
term.
   
   Not acceptable: a cap inside `MakeSequencedChanWithDiscard` that stops 
reading `source` once `pq.Len()` reaches a limit. The reorder goroutine is the 
only reader of `records`, which has capacity `numWorkers`. Once it stops 
reading, the other workers fill `records` with batches that cannot be emitted, 
the lagging worker's send blocks behind them, and the heap can never drain 
because the batch it is waiting for is stuck in the channel. Every goroutine 
waits on another and the scan hangs. The credit design does not have this 
failure because the reorder goroutine never stops reading; backpressure is 
applied to producers before they send.
   
   Fallback: if the credits regress `BenchmarkArrowScanManyFilesAndBatches` 
(table/arrow_scanner_bench_test.go:99) or cannot be made clean under `-race`, 
the PR instead narrows the doc at table/rewrite_data_files.go:258-266 to say 
the bound holds only with `WithCompactionScanConcurrency(1)`, and says so in 
its description.
   
   ### Correctness requirements
   
   - Output order is unchanged: task index, then record index, as the 
comparators at table/arrow_scanner.go:2130-2156 define it. Error records still 
sort ahead of data (:2137, :2148).
   - Every batch is released exactly once on every exit path: consumed by the 
caller, discarded at close (table/internal/utils.go:99-105), drained by the 
iterator's deferred loop (table/arrow_scanner.go:2166-2170), or cut by the row 
limit.
   - A worker waiting for a credit wakes on `scanCtx.Done()`, so after an error 
or a caller that stops iterating, `wg.Wait()` at table/arrow_scanner.go:2309 
returns and `records` is closed at :2310. No goroutine outlives the iterator.
   - A worker waiting for a credit holds no channel slot and no lock, so a 
worker with free credits can always send.
   - `WithMaxConcurrency(1)` behaves as today: one worker, credits never 
contend, output identical.
   - No public API change. K is an internal constant chosen from the benchmark 
below. A public option is added only if the benchmark shows that no single 
default is acceptable.
   
   ### Validation plan
   
   - A benchmark in table/ with the gated IO described above, reporting the 
HeapInuse delta while task 0 is gated for M = 8, 32, 128 at concurrency 4. 
Before the change it reproduces the table above. After the change the delta is 
flat in M and below `numWorkers x K x 536.5 KB` plus a fixed overhead.
   - A unit test with a lagging task 0 that asserts the reorder heap never 
holds more than `numWorkers x K` entries, and a test that the lagging worker 
still produces once released while every other worker is out of credits.
   - A test that cancels the scan while workers are blocked on credits and 
checks that the iterator returns and every worker goroutine exits.
   - `BenchmarkArrowScanManyFilesAndBatches` 
(table/arrow_scanner_bench_test.go:99) before and after at several K with 
`-count=10`, compared with benchstat. The default K is the smallest value whose 
benchstat delta on that benchmark is within noise.
   - `go test -race` over the scanner, rolling, fanout and clustered writer 
suites, the same set the #1993 review ran.
   - The doc at table/rewrite_data_files.go:258-266 rewritten to state the 
bound in batches. A follow-up that runs compaction groups concurrently will 
multiply this bound by the number of groups in flight, so its sizing docs 
depend on this bound being real.
   
   ### Related
   
   - #1993 added `WithCompactionArrowBatchSize`, 
`WithCompactionRecordBatchBufferSize` and `WithCompactionParquetRowGroupLimit`; 
its approval review recorded this gap.
   - #1982 is the request for a compaction memory bound that #1993 answered; 
this issue is the part of that bound that does not hold yet.
   - #1970 shares decoded manifest lists across concurrent scans of one table. 
The follow-up that runs compaction groups concurrently builds on both #1970 and 
this bound.
   
   ### 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