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]