laskoviymishka commented on code in PR #2039:
URL: https://github.com/apache/iceberg-go/pull/2039#discussion_r4075697070


##########
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:
   These new primitives carry the whole correctness story of the PR, but they 
read without any of the invariant comments the rest of this file leans on 
(lazyPositionDeleteLoader, posDeleteAccumulator and friends all have them). Two 
invariants here aren't visible from the code at all: a credit rides only the 
Last record, and a failed attempt never releases its credit explicitly. That 
second one is safe only because the worker cancels and returns right after, 
abandoning the token; a future retry-instead-of-return would silently 
reintroduce a leak with nothing here to warn against it.
   
   While we're in here, send/fail aren't nil-safe the way acquire/release are. 
A zero-value recordSink{} no-ops through reserve() and then blocks forever on 
the nil-channel send, so I'd note that these must be built via newRecordSink, 
or have send/fail treat a nil out as a programmer error and panic loudly rather 
than hang.
   
   Two cosmetics too: acquire's `if c == nil` guard can't actually fire today 
since newRecordSink always sets credits, so either drop it or mark it defensive 
the way releasePerFilePosDeletes does; and the credits field references 
taskCredits a few lines before that type is declared, which reads backwards. 
wdyt?



##########
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:
   fail() is meant to be the single failure path, but the three delete-loader 
error branches down in recordBatchesFromTasksAndDeletes (around line 2321) skip 
it and push enumeratedRecord{Task, Err} straight onto records.
   
   The net effect matches today, since neither fail() nor those direct sends 
attach a credit, but that's now the same invariant maintained in two places and 
easy to miss. If fail() ever grows bookkeeping (metrics, tracing, a second 
credit type), those three sites silently won't get it.
   
   I'd route them through sink.fail(task, err) so there's one failure path to 
reason about. wdyt?



##########
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:
   Small thing, non-blocking: this gate fires for every worker count, including 
numWorkers == 1, where the runaway heap can't happen anyway since a single 
worker already runs its tasks strictly in sequence.
   
   The only cost there is a little legitimate read-ahead, since the worker now 
waits for the top-level consumer to dequeue the last batch before starting its 
next task, which is a strictly later point than before. Probably negligible, 
but if we wanted numWorkers == 1 to stay exactly as fast as before, we could 
pass a nil taskCredits sink in that case and skip reserve() entirely. Happy 
either way.



##########
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:
   This memory-bound doc is great for the compaction path, but #2038 was filed 
against the general scan (Table.Scan / GetRecords / ReadTasks), and that path 
has no equivalent caveat. Someone reading only the general scan docs would come 
away thinking memory is now file-count-independent, when a size-skewed group 
(one big file among many small ones) can still spike transiently if the big 
file decodes early and its Last batch waits behind a slow head task. A short 
cross-reference near GetRecords would close that gap.
   
   Separately, the formula here is a precise, operator-facing claim that 
nothing exercises: the two new tests cover generic Scan/ReadTasks, not a gated 
or lagging task driven through the clustered writer. I'd either soften it 
toward best-effort/approximate, or add a compaction-path test that drives the 
writer side so it can't silently drift from the code. Thoughts?



-- 
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