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]

Reply via email to