dabla commented on PR #62922:
URL: https://github.com/apache/airflow/pull/62922#issuecomment-5762925386

   ### Why the last five commits move iterated inputs to async
   
   While running the IterableOperator in production (Airflow 3.3.2, the 
operator applied as a monkey patch), one iterated task froze a few times a 
week: a batch of sub-tasks logged "Attempting running task" and nothing else, 
the task process showed zero context switches on every thread, and the 
supervisor kept heartbeating until the 24h execution timeout. No exception, no 
log line, nothing on the API server.
   
   A `faulthandler` dump of a frozen process (registered on `SIGUSR2`; 
`SIGUSR1` triggers Celery's inherited soft-time-limit handler) showed the 
deadlock:
   
   - `AsyncAwareExecutor.map` pulled the next sub-task's input on the **main 
thread between two `run_until_complete` calls**, i.e. with the event loop 
paused.
   - That pull was a synchronous supervisor call 
(`LazyXComSequence.__getitem__` -> `SUPERVISOR_COMMS.send`). With no loop 
running, `send()` takes the comms thread lock in blocking mode.
   - A sub-task coroutine parked mid-`asend` already held that lock and could 
only release it once the loop ran again. It never did.
   
   So the rule the SDK already enforces for a *running* loop (a sync call there 
raises `DeadlockImminentError`) had a blind spot: a sync call from the loop's 
thread while the loop is paused blocks instead of raising, and anything in 
flight on the loop is stuck behind it.
   
   **What changed, in commit order**
   
   1. *Add async reads to XComIterable and LazyXComSequence* - `aget` and 
`__aiter__` on both, going through `XCom.aget_one` / `asend`, so the two lazy 
XCom sources can be read from a coroutine.
   2. *Split PlainXComArg.resolve into map-index and pulled-value helpers* - no 
behaviour change, only so the next commit can share the logic.
   3. *Add XComArg.aresolve* - the async twin of `resolve`, abstract on the 
base. `PlainXComArg` pulls through `ti.axcom_pull` and computes its map indexes 
in a worker thread (counting upstream TIs is a blocking supervisor call); 
Map/Zip/Concat await their parts and reuse the sync validation.
   4. *Add aiter_values to expand inputs and XComArg* - the async twin of 
`iter_values`, abstract on `ExpandInput`, implemented by the four concrete 
inputs. XComArg sources go through `aresolve`; resolved values go through 
`aiterate`, which uses `async for` on async iterables, iterates in-memory 
containers in place, and advances any other iterator in a worker thread because 
its `next()` may be a blocking XCom read.
   5. *Pull iterated sub-task inputs on the event loop to stop the 
IterableOperator deadlock* - `AsyncAwareExecutor.map` now takes async iterables 
only and pulls and submits from a coroutine on the running loop; the `while 
pending` loop and result streaming are unchanged and nothing is materialised. 
`IterableOperator` feeds it an async generator over `expand_input.aiter_values`.
   
   The sync `map` was dropped rather than kept alongside: with the pull awaited 
on the loop, a sync SDK call in an input now meets the running-loop check and 
raises immediately instead of hanging, which is the behaviour we want.
   
   **Tests added for the regression**
   
   - Executor: items are pulled on the running loop and only when a worker slot 
frees; a lock held across an `await` by a mapped coroutine that the input 
source also needs completes instead of deadlocking (run in a thread with a join 
timeout).
   - XComIterable / LazyXComSequence: async iteration and `aget` go through 
`aget_one` / `asend` only; the blocking `send` is asserted untouched.
   - XComArg: `aresolve` on all four types, missing-XCom semantics matched 
against `resolve` case by case, map indexes computed off the loop thread, 
`aiter_values` with the sync `resolve` rigged to fail.
   - Expand inputs: `aiter_values` for all four types including laziness of the 
dict product and async-only sources; `aiterate` advancing plain iterators off 
the loop thread.
   - Operator: `execute` consumes `aiter_values` on the running loop with 
`iter_values` rigged to fail; an XComArg input resolved without the sync 
`resolve`.
   
   Verified against the production DAG in a local stack: the previously 
freezing task ran its 114 sub-tasks to completion, and a checkpointed retry 
resumed correctly on the async path.
   
   ---
   Drafted-by: Claude Code (Fable 5.1); reviewed by @dabla before posting
   


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

Reply via email to