dabla commented on code in PR #62922:
URL: https://github.com/apache/airflow/pull/62922#discussion_r4094020782
##########
task-sdk/src/airflow/sdk/definitions/xcom_arg.py:
##########
@@ -180,6 +205,15 @@ def concat(self, *others: XComArg) -> ConcatXComArg:
def resolve(self, context: Mapping[str, Any]) -> Any:
raise NotImplementedError()
+ async def aresolve(self, context: Mapping[str, Any]) -> Any:
Review Comment:
Fixed in 01c851d1d6. `aresolve` now has a working default on the base:
`resolve` run through `asyncio.to_thread`, safe for the same reason
`aiterate`'s fallback is, since a blocking supervisor call from a worker thread
waits for in-flight `asend` calls instead of deadlocking with them.
`iter_values` and `aiter_values` already had working defaults built on those
two, so a subclass implementing only `iter_references` and `resolve` now works
end to end on the iterated path; the class docstring says so under a
"Subclassing" heading, and a test with exactly such a subclass pins all three.
The built-in subclasses keep their `ti.axcom_pull` overrides as the fast path.
Your point was well taken: two of our own in-house `XComArg` subclasses would
have hit exactly this.
---
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]