KranzL commented on issue #2038: URL: https://github.com/apache/iceberg-go/issues/2038#issuecomment-5777711499
PR #2039 bounds the heap with one task credit per worker instead of the K batches per worker proposed above. Each worker goroutine in `recordBatchesFromTasksAndDeletes` owns a `taskCredits` channel of capacity `maxInFlightTasksPerWorker` (table/arrow_scanner.go:1251 on the PR branch, value 1). It acquires the credit before it takes a task from `taskChan` (:2300), selecting on the scan context while it waits (:1259-1270), and `recordSink.send` (:1291) attaches the credit to the task's `Last` record. The credit is released when the consumer takes that record (:2247), when the iterator's deferred drain takes it (:2220), or when `MakeSequencedChanWithDiscard` discards it at close (:2209). A worker therefore cannot start its next task while any batch of its previous task is in the records channel, the reorder heap or the sequenced channel. While the head task lags, each free worker finishes one task and parks, so the heap holds at most `numWorkers - 1` tasks. table/internal/utils.go is unchanged. I built the per-batch version first, with the credit reserved before the worker decodes each batch and released when the consumer takes it, and measured it on `BenchmarkArrowScanManyFilesAndBatches` (32 files of 32 one-row batches, concurrency 4). Apple M3 Pro, 11 cores, go1.25.9 darwin/arm64, extra_properties_0 median ns/op. These runs were not interleaved with the baseline, so they carry up to 8% of thermal drift; the differences are larger than that. | credits per worker | median | runs | against de53d44 | |---|---:|---:|---:| | de53d44 (main plus the benchmark commit) | 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 is allowed to run ahead are what gives the scan its 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 a bound in batches is no tighter than one task per worker. The per-task credit 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. With the per-task credit the same benchmark is within noise of main (geomean +0.57%, n = 10 per cell, interleaved binaries, benchstat), and the lagging-task benchmark's peak HeapInuse delta at M = 128 goes from 94.6 to 96.6 MB down to 2.28 to 2.34 MB over 5 runs. The full tables and the rewritten doc for `WithCompactionArrowBatchSize` are in the PR description. -- 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]
