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]
