|
This AIP is a proposition to support an alternative way of expanding multiple XCom’s on operators within the need to know how may tasks will be expanded in advance, which is the case with streams or iterables.
That’s why I would also like to introduce then notion of streamable XCom's. In the current Airflow implementation, only list and dicts are supported as XCom collection types on which can be expanded on (e.g. SchedulerDictOfListsExpandInput and SchedulerListOfDictsExpandInput).
I would like to add support for iterable's (e.g. lazy evaluation collection). The reason why I would like to introduce the support of iterables is to allow the implementation of streamable XCom's. In case of the HttpOperator or the MSGraphAsyncOperator for example, when those return pages results, the operator depending on the results of that operator needs to wait until all pages have been loaded before being able to process them.
This has 4 disadvantages:
Is there anything special to consider about this AIP? Downsides? Difficulty in implementation or rollout etc?
You could argue that you could write a PythonOperator or a task decorated method in which you loop over the multiple inputs from the XCom to pass as an argument to the operator or even a hook. While the later would be a valuable solution, the first one wouldn’t as it’s a bad practise to execute an operator from within a PythonOperator, see the discussion about this topic on the devlist.
There is already an safeguard implemented for this which checks if an operator is executed from a PythonOperator, and if so, logs a warning stating an operator cannot be called outside of a TaskInstance. In the future this will probably become prohibited and will raise an AirflowException in that case.
But let’s hypothetical assume this would still be allowed, how would you loop the inputs from an XCom to an operator if that operator is deferrable? You would need to take multiple aspects into account.
First of all, you would need to catch the raised TaskDeferred exception, as this is how a deferrable operator works. The TaskDeferred exception contains the triggerer associated to the deferred operator to be executed, which behinds the scenes returns an async generator. This means that you will have to cope with the event loop of asyncio to be able to run the async triggerer from your PythonOperator, as the later one isn’t executed in an async way. Maybe this is also a good moment to start thinking of natively supporting async method’s in the PythonOperator without worrying about coping with the event loop (e.g. PythonTriggerer)?
Next to that, once you achieved to execute the triggerer, you will also have to check if a next_method was specified, which has to be executed on the deferred operator once the trigger has completed.
And last but not least, the execution of the next_method could also re-raise a TaskDeferred exception if the deferred operator implements the producer/consumer pattern, which means you’ll have to take into account recursion. For example the MSGraphAsyncOperator implements the producer/consumer pattern in such a way that the worker triggers the request to the MS Graph API, but instead of blocking the worker waiting for the response to arrive, releases it and delegates it to the triggerer, avoiding blocking workers unnecessarily. That way when the triggerer receives the response, it gives the received response back to the operator (e.g. worker) without blocking the worker while awaiting for the response.
That’s already a lot of technical challenges you have to solve if you want to execute a (deferrable) operator from within a loop in a PythonOperator.
Beside that, you maybe would also like to introduce some multithreading to speed up the processing instead of just looping in a sequential manner, unless you want sequential execution in the given input order, which is also a issue raised in this discussion. But there you will also have to be careful, because you just can’t execute a deferable operator having async code through a ThreadPoolExecutor.
Faster execution of multiple task instances related to one and the same operator within the same worker instance, having the advantage of sharing the same memory without overloading the Airflow scheduler with multiple task instances related to one and the same operator. Avoiding unnecessary waste of time by waiting until a paged result is completely loaded by an operator before passing it to the next operator.
![Airflow > [WIP] AIP-88: Streaming & lazy task expansion via long-running triggers > image_720.png](/confluence/download/attachments/334760511/image_720.png?version=1&modificationDate=1749032267000&api=v2)
In the screenshot above, you can see the new lazy expandable task mapping implementation in Airflow 3. It demonstrates how the MSGraphAsyncOperator returns a deferred iterable (as an XCom), allowing downstream tasks like SQLInsertRowsOperator to expand lazily as the scheduler (e.g. TaskMap) iterates over the paged results.
This means the SQLInsertRowsOperator does not have to wait for all pages to be fetched before expanding into mapped tasks. Instead, the MSGraphAsyncOperator initially returns only the first page, and as the scheduler expands the mapped tasks, it fetches additional pages on demand.
This approach offers significant improvements over the traditional model, where the MSGraphAsyncOperator would have to fetch all pages upfront before expansion could begin—resulting in slower performance and higher memory usage.
Simplifies concurrent or sequential iteration of multiple XCom's over one (deferrable) operator within the same task instance.
With the current proposal, task expansion is entirely managed by the scheduler, rather than by the MappedOperator itself. As a result, XComs are now resolved within the scheduler, not inside the operator. The primary benefit of this approach is that the length of the XCom no longer needs to be known upfront, enabling streaming-like, partial task expansion.
Initially, the scheduler expands only the first task instance from the MappedOperator. The remaining task instances are expanded asynchronously by the TaskExpansionJobRunner, allowing new tasks to be created incrementally as more data becomes available. Meanwhile, the executor can immediately begin running the already-expanded tasks.
This results in a more dynamic and efficient execution model, particularly valuable for deferred or paginated iterables, where fetching all data upfront would be slow or memory-intensive.
None
None
None, as this is just an alternative way of expanding XCom's on a operator, instead of calling the existing partial method of the MappedOperator, you can now call the iterate method.