KranzL commented on code in PR #2039:
URL: https://github.com/apache/iceberg-go/pull/2039#discussion_r4077479754
##########
table/arrow_scanner.go:
##########
@@ -1242,9 +1242,61 @@ func (as *arrowScan) addTaskProjectedFieldIDs(invariants
*arrowScanInvariants, t
}
type enumeratedRecord struct {
- Record tblutils.Enumerated[arrow.RecordBatch]
- Task tblutils.Enumerated[FileScanTask]
- Err error
+ Record tblutils.Enumerated[arrow.RecordBatch]
+ Task tblutils.Enumerated[FileScanTask]
+ Err error
+ credits taskCredits
+}
+
+const maxInFlightTasksPerWorker = 1
+
+type taskCredits chan struct{}
Review Comment:
Reworked in 149c395. The invariant you named is gone rather than documented:
reserve marks the credit as held, and the sink hands it to exactly one record,
the task's Last record in send or, when no Last record was sent, the error
record in fail. processRecords can send Last and then still return an error
from the reader or the deferred close, so fail hands off nil in that case and
the consumer releases each credit once. A worker that failed and went back to
reserve would block only until its error record leaves the sequenced channel,
which the iterator's drain loop guarantees, so a retry-and-continue path would
neither leak a credit nor hang. TestRecordSinkHandsOffCreditOnce
(table/arrow_scanner_reorder_test.go) covers fail without a Last record, Last
followed by fail, and that reserve blocks while the credit is out.
recordSink is a pointer now and is built only by newRecordSink (the two
worker loops and the three pruning tests), so a nil *recordSink panics on first
use instead of hanging on a nil channel. The nil guard in acquire is dropped;
the one in release stays because every record other than Last and error records
carries a nil credit. taskCredits is declared before enumeratedRecord.
I kept the primitives free of comments and put the invariant into code and
tests instead.
##########
table/arrow_scanner.go:
##########
@@ -1242,9 +1242,61 @@ func (as *arrowScan) addTaskProjectedFieldIDs(invariants
*arrowScanInvariants, t
}
type enumeratedRecord struct {
- Record tblutils.Enumerated[arrow.RecordBatch]
- Task tblutils.Enumerated[FileScanTask]
- Err error
+ Record tblutils.Enumerated[arrow.RecordBatch]
+ Task tblutils.Enumerated[FileScanTask]
+ Err error
+ credits taskCredits
+}
+
+const maxInFlightTasksPerWorker = 1
+
+type taskCredits chan struct{}
+
+func newTaskCredits() taskCredits {
+ return make(taskCredits, maxInFlightTasksPerWorker)
+}
+
+func (c taskCredits) acquire(ctx context.Context) error {
+ if c == nil {
+ return nil
+ }
+
+ select {
+ case c <- struct{}{}:
+ return nil
+ case <-ctx.Done():
+ return context.Cause(ctx)
+ }
+}
+
+func (c taskCredits) release() {
+ if c != nil {
+ <-c
+ }
+}
+
+type recordSink struct {
+ out chan<- enumeratedRecord
+ credits taskCredits
+}
+
+func newRecordSink(out chan<- enumeratedRecord) recordSink {
+ return recordSink{out: out, credits: newTaskCredits()}
+}
+
+func (s recordSink) reserve(ctx context.Context) error {
+ return s.credits.acquire(ctx)
+}
+
+func (s recordSink) send(rec enumeratedRecord) {
+ if rec.Record.Last {
+ rec.credits = s.credits
+ }
+ s.out <- rec
+}
+
+func (s recordSink) fail(task tblutils.Enumerated[FileScanTask], err error) {
Review Comment:
Done in 149c395: the three loader error sites call sink.fail(task, err) and
then cancel and return, so fail is the only path that puts an error record into
records. fail is a plain send, as the recordsFromTask path already was. The
sequencing goroutine in MakeSequencedChanWithDiscard reads records until every
worker has exited, and the iterator's deferred loop drains the sequenced
channel after the caller stops, so the select on scanCtx.Done() at those sites
was not what kept them from hanging. With the change in the other thread, fail
also carries the worker's credit when no Last record was sent.
##########
table/arrow_scanner.go:
##########
@@ -2240,7 +2295,11 @@ func (as *arrowScan)
recordBatchesFromTasksAndDeletes(ctx context.Context, tasks
for range numWorkers {
go func() {
defer wg.Done()
+ sink := newRecordSink(records)
for {
+ if err := sink.reserve(scanCtx); err !=
nil {
Review Comment:
Measured instead of special-cased. BenchmarkArrowScanManyFilesAndBatches
with concurrency set to 1 in the benchmark (32 files of 32 one-row batches, so
1024 task boundaries per op), test binaries for de53d44 and 149c395 run
interleaved, 5 rounds with -test.count=2 each, benchstat over n = 10 per cell,
Apple M3 Pro, go1.25.9:
| sub-benchmark | de53d44, sec/op | 149c395, sec/op | change |
|---|---:|---:|---|
| extra_properties_0 | 286.2m ± 18% | 285.9m ± 19% | ~ (p=0.684) |
| extra_properties_16 | 282.9m ± 10% | 300.4m ± 21% | ~ (p=0.280) |
| extra_properties_64 | 281.9m ± 8% | 288.9m ± 17% | ~ (p=0.481) |
| geomean | 283.6m | 291.6m | +2.82% |
No sub-benchmark differs at p < 0.05. The machine was on battery power for
these runs, so the absolute times are about 2.2x the ones in the description;
the same is true of the concurrency-4 numbers I posted in the top-level
comment, which were taken in the same session.
I kept the gate for every worker count because the bound on
WithCompactionArrowBatchSize assumes it: without the gate a single worker can
leave its previous task's last batch in the records channel and the sequenced
channel while it decodes the next task, which would add two batches to the
formula only when workers is one.
##########
table/rewrite_data_files.go:
##########
@@ -257,13 +261,25 @@ func WithCompactionScanConcurrency(n int)
CompactionGroupOption {
// WithCompactionArrowBatchSize caps the number of rows decoded per
// Arrow record batch while reading the group's tasks, forwarded to the
-// scan as [WithArrowBatchSize]. Together with
-// [WithCompactionRecordBatchBufferSize] it bounds the memory held by
-// the record pipeline specifically: buffered batches times rows per
-// batch. Delete-side memory is not covered — positional deletes and
-// deletion-vector bitmaps for the group's tasks are materialized up
-// front and sized by delete volume, not by these knobs. A non-positive
-// value keeps the table's read.parquet.batch-size property.
+// scan as [WithArrowBatchSize]. The record pipeline holds at most
Review Comment:
Both in 149c395. Scan.ReadTasks (table/scanner.go) now states the bound for
the general scan: each worker holds the decoded batches of at most one task
until the iterator has returned them, so the batches waiting behind a slow
early task total at most the worker count times the rows in the largest task,
rounded up to a multiple of the batch size, and one large file among small ones
still holds its whole decoded contents while it waits. ToArrowRecords and
WithMaxConcurrency point at that paragraph.
For the formula, TestExecuteCompactionGroupRecordPipelineBounded
(table/arrow_scanner_reorder_test.go) drives ExecuteCompactionGroup with
WithCompactionScanConcurrency(4) and WithCompactionRecordBatchBufferSize(2)
over 32 single-batch files inside testing/synctest, gating task 0 in Open and
the output file in Create. Once every goroutine is durably blocked it checks
that exactly 3 tasks were fully read while task 0 was gated, (workers - 1) x
maxInFlightTasksPerWorker, then opens task 0 and checks that exactly 8 were
fully read while the writer was gated: workers x maxInFlightTasksPerWorker +
recordBatchBufferSize + 2, which are the two terms of the formula counted in
tasks. With the credit gate disabled locally the two counts are 31 and 32. A
change that adds a queue stage on either side moves the second count and fails
the test, so the formula text is unchanged and stays exact rather than
best-effort.
--
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]