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]
