You are viewing an old version of this page. View the current version.

Compare with Current View Page History

« Previous Version 48 Next »

Status

StateDraft
Discussion Thread

https://lists.apache.org/thread/zxpbhwb0nb28st4q9ddnf2900d828xsm

Vote Thread
Vote Result Thread
Progress Tracking (PR/GitHub Project/Issue Label)
Date Created

17.12.2024 07:36

Version Released
AuthorsDavid Blain

Short summary

Allow Triggerers to yield and have the scheduler/process handle multiple events in a single trigger run, enabling true streaming-style workflows for async operators (e.g., 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

  1. 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.

  2. Allow operators that know how to consume events to expand mapped tasks or XCom-driven iterables lazily as events arrive.

  3. Preserve serializability guarantees (no unserializable callables in trigger args).

  4. 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.

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 during run(). It still must be serializable. If the trigger performs pagination, it should do that inside its run() loop and yield each 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_method handler incrementally.

  • Operator / Worker:

    • Operators supporting streaming implement a handle_trigger_events(events: Sequence[TriggerEvent], context) callback (or extend next_method to 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 TaskDeferred to 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_trigger exists and is used (Airflow 3.0+).

API changes / additions

  1. Triggerer run behavior: no signature change, but scheduler processing semantics change: treat yield as a stream of events, not just the first.

  2. 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.

  3. Operator callback:

    • Add next_method(self, event: TriggerEvent | List[TriggerEvent], context) accept either single or batched events, or an optional handle_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.

  4. DAG parsing checks: if start_from_trigger is enabled and start_trigger_args contains a non-serializable callable, raise at parse time (existing check you already mentioned).

Ordering & consistency concerns

  • Use sequence_index to 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_args types 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_events or accept repeated next_method calls. Provide an adapter shim that calls an operator’s single-event next_method repeatedly 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 TriggerEvent envelope,

    • Ensure start_trigger_args are serializable,

    • Implement handle_trigger_events to 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_trigger but 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

  • No labels