dabla commented on code in PR #62922:
URL: https://github.com/apache/airflow/pull/62922#discussion_r4093443410
##########
task-sdk/src/airflow/sdk/definitions/_internal/expandinput.py:
##########
@@ -184,6 +363,85 @@ def iter_references(self) -> Iterable[tuple[Operator,
str]]:
if isinstance(x, XComArg):
yield from x.iter_references()
+ def iter_values(self, context: Mapping[str, Any]) -> Iterable[Any]:
+ from airflow.sdk.definitions.xcom_arg import XComArg
+
+ def _to_iterable(v: Any) -> Iterable:
Review Comment:
Fixed in f55fb22809. Both paths now go through a single module-level
`_to_iterable`, so they cannot drift again: a Mapping hands over its `.items()`
view on both sides (no more `list()` on the async side), any iterable, sync or
async, passes through untouched, and scalars become a one-element tuple.
`aiterate` now counts dict and mapping views as in-memory containers, so the
view is walked in place rather than one item per worker thread. Tests cover the
helper directly and pin that `iter_values` and `aiter_values` yield the same
combinations for a dict argument.
---
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]