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]

Reply via email to