DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
Status
| Page properties | ||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
Short summary
Allow Triggerers to yield multiple events in a single trigger run, enabling true streaming-style workflows for async operators
Motivation
This AIP was a proposition to support an alternative way of expanding multiple XCom’s on operators within the same task instance (e.g. worker instance), especially in Airflow 2.x, hence we will wait until Airflow 3.0 is released and re-evaluate how well it will then handle +10k XCom’s.
Beside that, I would like to introduce streamable XCom's. In the current Airflow implementation, only list, dicts and simple build-in Python types are supported as XCom types. I would like to add support for iterable's. 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 2 disadvantages:
- performance, as depending operators wait until all pages have been loaded before being able to start
- memory usage, as all pages have to be loaded into memory as an XCom before the next operator can start
Next to that, in Airflow 2.x we encountered performance issues when we had to expand like +10k XCom’s (e.g. which is also the case when processing pages results)
The performance issues where 2 fold:
- Even though we use the CeleryExecutor, pgbouncer as a connection pool to our postgres database and defined dedicated queues for DAG’s that had to process a lot of task instances, the Airflow UI became unresponsive due to the large amount of mapped task instances being created.
- A lot of time was spend orchestrating the execution state of all mapped task instances, even though we have a pool of workers ready to run. The problem here is the communication between the workers and the scheduler, which most of the time took more time than the task execution itself. Also we discovered that the scheduler was actually DDOS-ing the database, due to the large amount of mapped task instances.
Another issue with the current partial/expand mechanism is that it expects to know how many tasks will be expanded in advance (e.g. sequence), which is not the case with an iterable (and streaming of course).
, HttpOperator/MSGraphAsyncOperator). This reduces repeated context switching and "ping-pong" between scheduler → worker → triggerer, enabling efficient in-trigger pagination and lazy task expansion.
Motivation (why)
Many async operators implement pagination by deferring to triggers and then reentering the operator repeatedly. Even with
start_from_trigger, the current system processes only the first TriggerEvent yielded; additional pages require repeated defer/resume cycles.Repeated context switches create scheduler/worker/triggerer overhead and slow throughput for high-volume paginated APIs.
A streaming-capable triggerer can yield many events in one run; letting Airflow process those events (without full defer/resume for each page) produces substantial performance gains and allows lazy expansion of downstream tasks as pages arrive.
Goals
Allow triggers to yield many events during a single run and have those events processed incrementally by the operator runtime/scheduler without repeated operator deferral cycles.
Allow operators that know how to consume events to expand mapped tasks or XCom-driven iterables lazily as events arrive.
Preserve serializability guarantees (no unserializable callables in trigger args).
Be backward-compatible: existing triggers/operators continue to work.
Non-goals
Making arbitrary XComs iterable across the cluster. (That requires additional XCom semantics; out of scope.)
Allowing triggers to carry non-serializable callables/closures. Triggers must remain serializable.
Demo
| Widget Connector | ||
|---|---|---|
|
High-level design
Key idea
Extend the triggerer-to-scheduler event path so that a single trigger run may produce a stream of events which the scheduler can incrementally deliver to the operator context (worker/scheduler) and allow the operator to handle each event without a full defer/resume roundtrip per page.
Components & responsibilities
Triggerer: can yield multiple
TriggerEvents duringrun(). It still must be serializable. If the trigger performs pagination, it should do that inside itsrun()loop andyieldeach page as an event.Scheduler: must accept and process multiple events from the same trigger-run. Instead of ignoring all but the first yielded event, scheduler accepts the events and binds them to the operator’s
next_methodhandler incrementally.Operator / Worker:
Operators supporting streaming implement a
handle_trigger_events(events: Sequence[TriggerEvent], context)callback (or extendnext_methodto accept a batch/stream), which may:process events and optionally expand mapped tasks or create XComs,
decide whether trigger should continue producing (e.g., requests more pages) — but triggers are the ones that control pagination inside
run().
Operators still raise
TaskDeferredto start the trigger if they cannot proceed synchronously.
Start-from-trigger: still used — when present we skip the initial worker-run step. The AIP assumes
start_from_triggerexists and is used (Airflow 3.0+).
API changes / additions
Triggerer run behavior: no signature change, but scheduler processing semantics change: treat
yieldas a stream of events, not just the first.Scheduler/Triggerer protocol:
Add an event envelope with fields
{trigger_run_id, sequence_index, payload, is_last}for each yielded event. This helps ordering and safe incremental processing.
Operator callback:
Add
next_method(self, event: TriggerEvent | List[TriggerEvent], context)accept either single or batched events, or an optionalhandle_trigger_events(self, events, context)method.For backward compatibility, if operator implements the old
next_method(event, context), the scheduler will call it for each event individually.
DAG parsing checks: if
start_from_triggeris enabled andstart_trigger_argscontains a non-serializable callable, raise at parse time (existing check you already mentioned).
Ordering & consistency concerns
Use
sequence_indexto process events in order from the same trigger run.If operator processing of an event fails, the scheduler should:
surface the failure (normal task retry semantics) and
stop accepting further events for that run until retry/resume behavior is resolved.
If the trigger signals
is_last, scheduler can mark the trigger-run completed.
Security / serializability
Triggers and all trigger args must remain serializable (no lambdas/closures).
DAG parsing validates
start_trigger_argstypes and raises early if a non-serializable callable is present (your existing behavior).
Backward compatibility / migration
Existing triggers/operators continue to work: scheduler defaults to single-event processing if operator doesn’t opt-in.
Operators that want streaming implement the new batched
handle_trigger_eventsor accept repeatednext_methodcalls. Provide an adapter shim that calls an operator’s single-eventnext_methodrepeatedly if the operator hasn’t implemented streaming (so no breakage).Update common async operators (MSGraphAsyncOperator, HttpAsyncOperator) to implement in-trigger pagination and yield many
TriggerEvents.
Tests & acceptance criteria
Unit tests for:
multiple-event processing ordering,
failure during mid-stream event handling,
serializability check during DAG parse.
Integration tests:
MSGraphAsyncOperator with 10+ pages yields vs old approach measured in scheduler/worker/triggerer traces (assert fewer context-switch cycles).
Performance benchmarks demonstrating reduced scheduler <-> worker <-> triggerer hops.
Example operator migration notes
Refactor operator to:
Move pagination into the trigger
run()loop,Yield each page as a
TriggerEventenvelope,Ensure
start_trigger_argsare serializable,Implement
handle_trigger_eventsto consume events, create XComs, and expand tasks lazily.
Open questions / future work
Native iterable XComs (separate AIP).
Fine-grained backpressure between triggerer and scheduler (if extremely high event rates).
Event batching strategies (time / size) to trade off latency vs. scheduler invocation overhead.
Sequence diagrams
Below are two diagrams:
- Current/present behavior (with
start_from_triggerbut only first yielded event processed leading to repeated defer/resume cycles). - Proposed behavior (trigger yields many events in one run; scheduler forwards them incrementally to the operator; fewer context switches).
vs
Proposal
Possible workaround?
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.
What problem does it solve?
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.
Why is it needed?
Simplifies concurrent or sequential iteration of multiple XCom's over one (deferrable) operator within the same task instance.
Are there any downsides to this change?
With the current proposal the iteration logic is handled within the same IteratorOperator instead of the scheduler, but the current proposition could be seen as a facilitation to iterate multiple XCom's over one operator. This also mean you won't see the progres of each individual processed XCom with the Airflow UI, until it completely done. Another downside in the current proposition is that there isn't any state persisted regarding the progress of the already processed XCom's within the iteration, so there we should then have to think about a solution to persist the state of already iterated XCom's, as we can't persist it as an XCom as those are flushed each time a task fails for idempotency reasons. The better solution would be to have this alternative way of expanding task instances within the scheduler itself.
Which users are affected by the change?
None
How are users affected by the change? (e.g. DB upgrade required?)
None
What is the level of migration effort (manual and automated) needed for the users to adapt to the breaking changes? (especially in context of Airflow 3)
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.
Other considerations?What defines this AIP as "done"?






