KranzL opened a new pull request, #2039: URL: https://github.com/apache/iceberg-go/pull/2039
Fixes #2038. ### What changed Each scan worker now holds a task credit. `recordBatchesFromTasksAndDeletes` (table/arrow_scanner.go:2285) creates one `recordSink` per worker goroutine (:2298) with a `taskCredits` channel of capacity `maxInFlightTasksPerWorker` (:1251, value 1). The worker acquires the credit before it takes a task from `taskChan` (:2300) and selects on the scan context while waiting (:1259-1270), so cancellation still stops every worker. `recordSink.send` (:1291) attaches the credit to the task's `Last` record. The credit is released when the consumer in `createIteratorWithCleanup` takes that record (:2247), when the iterator's deferred drain loop takes it (:2220), or when `MakeSequencedChanWithDiscard` discards it at close through the existing discard callback (:2209). Error records carry no credit. The position delete producer in `makePositionDeleteRecordsForFilter` (table/transaction.go:3233-3235) gets the same sink because it feeds the same `createIterator`. `MakeSequencedChan`, `MakeSequen cedChanWithDiscard` and table/internal/utils.go are unchanged. A worker therefore cannot start task j while any batch of its previous task is still in the records channel, the reorder heap or the sequenced channel. While the head task lags, each of the other workers finishes one task and parks, so the heap holds at most `numWorkers - 1` tasks instead of every task the free workers can reach. Decoding inside a task is not throttled, and the worker on the head task never waits. ### Why a credit per task and not per batch #2038 proposed K batches of credit per worker. I implemented that first, with the credit reserved before the worker decodes its next batch and released when the consumer takes the batch, and measured it on `BenchmarkArrowScanManyFilesAndBatches` (32 files of 32 one-row batches, concurrency 4). These runs were not interleaved with the baseline, so they carry the machine's thermal drift, which I measured at up to 8% between two baseline runs; the differences below are larger than that. | credits per worker | extra_properties_0, ns/op median | runs | against main | |---|---:|---:|---:| | main at de53d44 | 69.4 ms | 5 | | | K = 2 | 141.6 ms | 5 | +104% | | K = 8 | 127.9 ms | 3 | +84% | | K = 32 | 78.7 ms | 3 | +13% | The batches a free worker may run ahead are the scan's parallelism. With fewer credits than batches per file, the three free workers stop after K batches and the head worker decodes alone, so the scan runs close to serially. K has to exceed the batch count of the largest file to avoid that, and at that point the bound in batches is no tighter than one task per worker. A credit per task keeps every worker decoding its own file at full speed and bounds the heap by task count, which for a compaction group is the worker count times the largest input file. ### Measurements Apple M3 Pro, 11 cores, go1.25.9 darwin/arm64. Before is de53d44 (origin/main 833e10c plus this branch's benchmark and test commit), after is this branch. `go test ./table/ -run xxx -bench BenchmarkArrowScanReorderHeapLaggingTask -benchmem -count=5`. gated-batches is the number of non-head files fully read while task 0 was blocked in Open. peak-heap-delta-MB is the largest HeapInuse increase over the pre-scan baseline while the gate was closed. | files M | gated-batches before | gated-batches after | peak-heap-delta-MB before, 5 runs | peak-heap-delta-MB after, 5 runs | |---:|---:|---:|---|---| | 8 | 7 | 3 | 4.66, 4.70, 4.68, 4.68, 4.63 | 2.05, 2.03, 2.05, 2.01, 2.06 | | 32 | 31 | 3 | 24.59, 20.05, 25.11, 23.82, 20.34 | 2.11, 2.19, 2.09, 2.24, 2.07 | | 128 | 127 | 3 | 95.09, 94.61, 95.67, 96.55, 94.86 | 2.28, 2.31, 2.29, 2.28, 2.34 | The ns/op column of that benchmark is not a throughput number on either side: the benchmark opens the gate after 100 ms without a file close, which on the fixed code is how it detects that the free workers have parked. `BenchmarkArrowScanManyFilesAndBatches`, test binaries for both commits run interleaved (5 rounds of before, after with `-test.count=2` each, so n = 10 per cell), compared with benchstat: | sub-benchmark | before, sec/op | after, sec/op | change | |---|---:|---:|---| | extra_properties_0 | 75.62m ± 1% | 75.94m ± 1% | ~ (p=0.089) | | extra_properties_16 | 75.94m ± 2% | 75.77m ± 2% | ~ (p=0.684) | | extra_properties_64 | 75.07m ± 1% | 76.22m ± 1% | +1.52% (p=0.002) | | geomean | 75.54m | 75.98m | +0.57% | B/op, allocs/op, batches/op and files/op are unchanged. The same loop also ran a build with two tasks of credit per worker: -0.94% geomean, every sub-benchmark within noise. One task per worker is the tighter bound and is within the same noise, so that is the default. ### Docs table/rewrite_data_files.go:249-262 (`WithCompactionScanConcurrency`) now says the scan's worker count multiplies the read-side term of the bound. table/rewrite_data_files.go:264-283 (`WithCompactionArrowBatchSize`) states the bound as `workers x (rows in the largest task + n) + (recordBatchBufferSize + 2) x n` rows, explains each term, and keeps the list of what it does not cover: the Parquet reader's own buffers and delete-side memory. Before this change the read-side term had no bound above scan concurrency one, so n did not bound the pipeline at the default concurrency. ### Tests `TestArrowScanReorderHeapBoundedWhileHeadTaskLags` (table/arrow_scanner_reorder_test.go:194) scans 32 single-batch files at concurrency 4 with task 0 gated in Open, inside testing/synctest so the count is read once every scan goroutine is durably blocked. On de53d44 it fails with "scan held 31 batches from the 31 out-of-order tasks while task 0 was gated, ... bound is numWorkers x K = 4 x 2 = 8". Here it passes with 3 completed out-of-order tasks against a bound of `(numWorkers - 1) x maxInFlightTasksPerWorker = 3`, then opens the gate and checks that all 8192 rows arrive and that the CheckedAllocator is empty. `TestArrowScanReorderCancelWhileWorkersWaitForTaskCredits` (table/arrow_scanner_reorder_test.go:248) uses the same setup, cancels the context while the three free workers are parked on their credits, checks that the iterator returns context.Canceled, opens the gate, and checks that every scan goroutine exits (synctest fails the test otherwise) and that every batch was released. `TestProcessRecordsUsesRowGroupFilterForPruning` and its two siblings in table/arrow_scanner_pruning_internal_test.go call `processRecords` directly and now wrap their channel in `newRecordSink`. ### Validation ``` go test ./table/ -run 'TestArrowScanReorder' -race -count=5 go test ./table/ -race -count=1 -run 'Rolling|Fanout|Clustered|ArrowScan|Reorder|RewriteDataFiles|ExecuteCompactionGroup' go test ./table/... -race -count=1 make lint ``` -- 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]
