Skip to main content

Module coop_amplifiers

Module coop_amplifiers 

Source
Expand description

Cooperative yielding for input-amplifying operators (#217).

DataFusion’s EnsureCooperative instruments LEAF streams only: budget is consumed per batch a leaf produces. An operator that amplifies its input — a cross or nested-loop join whose output is orders of magnitude larger than its input, or an unnest — drains its tiny budget-aware inputs in microseconds and then computes budget-free: a 5-way cross join over five 100-row VALUES tables feeds an aggregate 10^10 rows while consuming 5 units of budget, so its poll never yields and no timeout, cancel watcher, or select! arm can ever run (measured: a 2 s tokio::time::timeout armed around it did not fire in 7+ minutes).

The fix is one wrapper: put a CooperativeExec on top of each amplifier so budget is also consumed per OUTPUT batch. The stream then returns Pending every ~128 batches (~1M rows), which is what makes the executor’s cancel watcher and every timeout real for this operator class. datafusion-proto round-trips CooperativeExec, so distributed fragment encoding is unaffected.

Structs§

CooperativeAmplifiers
Wraps input-amplifying operators in CooperativeExec so their output participates in cooperative scheduling. See the module docs for why the default leaf-only instrumentation is not enough.