dabla commented on code in PR #62922: URL: https://github.com/apache/airflow/pull/62922#discussion_r4093944743
########## task-sdk/src/airflow/sdk/definitions/_internal/expandinput.py: ########## @@ -41,6 +51,32 @@ OperatorExpandKwargsArgument = Union["XComArg", Sequence[Union["XComArg", Mapping[str, Any]]]] +async def aiterate(iterable: Any) -> AsyncIterator[Any]: + """ + Iterate ``iterable`` from a coroutine without a blocking SDK call on the loop thread. + + An async iterable (``XComIterable``, ``LazyXComSequence``) is consumed with ``async for``, so + its reads go through ``asend``. In-memory containers are iterated in place. Anything else may + fetch on ``next()`` through a synchronous supervisor call, so each ``next()`` runs in a worker + thread: from there a blocking send waits for in-flight ``asend`` calls instead of deadlocking + with them (see ``AsyncAwareExecutor.map``). + """ + if hasattr(iterable, "__aiter__"): + async for item in iterable: + yield item + return + + if isinstance(iterable, (list, tuple, set, frozenset, range, deque)): + for item in iterable: + yield item + return + + iterator = iter(iterable) + sentinel = object() + while (item := await asyncio.to_thread(next, iterator, sentinel)) is not sentinel: Review Comment: Documented in ad236efea9, right below the reason: one `asyncio.to_thread` dispatch per item, so a 1000-item generator means 1000 thread hand-offs, and which sources avoid it (in-memory containers, now including dict and its views, are walked in place; an async iterable is consumed on the loop). --- Drafted-by: Claude 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]
