Skip to content

AIP-104: Task Iteration - #62922

Open
dabla wants to merge 27 commits into
apache:mainfrom
dabla:feature/dynamic-task-iteration
Open

dabla wants to merge 27 commits into
apache:mainfrom
dabla:feature/dynamic-task-iteration

Conversation

@dabla

@dabla dabla commented Mar 5, 2026 •

Copy link
Copy Markdown
Contributor

Was generative AI tooling used to co-author this PR?
  • [ x ] Yes (please specify the tool below)

Claude Code (Fable 5.1).

Description

This PR is the initial implementation of Iterable Tasks (IT), as discussed in the devlist and building upon the foundations of AIP-104. (Originally prototyped as "Dynamic Task Iteration"; renamed to Iterable Tasks following review feedback to avoid confusion with Dynamic Task Mapping.)

For further context on the use cases and performance benefits of IT, see this Medium Article and the new dynamic-task-mapping-vs-iteration.rst doc added in this PR, which compares IT with Dynamic Task Mapping (DTM) and Dynamic Task Batching in depth.

The XCom Database Constraint Challenge

While porting our internal "monkey-patched" version of IT (used since Airflow 2.x) to the core, I've identified a significant technical hurdle regarding XCom handling.

Around Airflow 2.10/2.11, a change was introduced to the database constraints for the XCom table. Specifically:

  • Current State: The DB prevents creating indexed XComs (map_index >= 0) unless a corresponding mapped TaskInstance exists in the task_instance table.
  • The Conflict: IT is designed to process multiple indexed XComs within a single Task Instance. Because there is no 1-to-1 mapping of a sub-task index to a physical TI row, the DB constraint blocks the insertion of these results.

The drawback is that XComs wouldn't automatically be removed from the database when a TaskInstance is deleted, which is the purpose of that constraint. So appending the index to the XCom key would be a good enough solution for IT, but not for DTM.

Current implementation in this PR

  • XComIterable (airflow.sdk.bases.xcom) appends the sub-task index directly to the XCom key (return_value_<index>) to bypass the constraint, and exposes the results as a lazy Sequence (__len__/__getitem__/__iter__) so a downstream task can consume them the same way it would consume an .expand() result.
  • XComIterable.flatten() returns a FlattenedXComIterable that lazily expands nested iterables (e.g. a list of pages into a single stream of items) without ever materializing the flattened stream in memory — __len__/__getitem__/__iter__ all speak consistently in flattened items.
  • Progress and crash recovery no longer rely on XCom alone. IterableOperator tracks per-sub-task progress in the AIP-103 Task State Store rather than XCom, and participates in Airflow's standard retry mechanism: if the IterableOperator's task instance is retried (or manually cleared before it finishes), already-succeeded sub-tasks are skipped and only pending/failed ones re-run. XCom is only (re)written per index once a sub-task succeeds; once every index has succeeded, checkpoints are dropped so a subsequent manual clear reruns all indices from scratch rather than replaying stale results.
  • Outlet asset/inlet events and on_kill propagation are handled per sub-task: a checkpointed sub-task replays its recorded outlet events on the following attempt instead of losing them, and killing the IterableOperator's task instance propagates to any sub-tasks still in flight.

I believe the cleanest long-term path is still to add a dedicated route in the Execution API that retrieves multiple XComs for a single TaskInstance by a list of keys in one round trip, so XComIterable.__getitem__/slicing don't need one request per element. I have a PR open to address this, intentionally split out of this PR.

This was also discussed in the devcall, see 2026-06-04 Dev Call Minutes.

AIP-104 itself was discussed again in the latest devcall, where the concerns raised there have also been addressed: 2026-09-10 Dev call Minutes.

Examples

The examples below assume an HTTP connection named pokeapi pointing to https://pokeapi.co.

Task Iteration

This example fetches a list of Pokémon from the PokéAPI and then uses Iterable Tasks (IT) to retrieve the details of each Pokémon. A single task instance processes all Pokémon URLs.

from airflow.sdk import dag, task
from airflow.providers.http.hooks.http import HttpHook, HttpAsyncHook

from pendulum import datetime

@dag(
    start_date=datetime(2025, 1, 1),
    schedule=None,
    catchup=False,
)
def pokemon_iteration():
    @task
    def list_pokemon() -> list[str]:
        response = HttpHook(
            http_conn_id="pokeapi",
            method="GET",
        ).run(
            endpoint="api/v2/pokemon?limit=100",
        )

        return [
            pokemon["url"].replace("https://pokeapi.co/", "")
            for pokemon in response.json()["results"]
        ]

    @task(
        retries=3,
        task_concurrency=2,
        show_return_value_in_logs=False,
    )
    async def get_pokemon(url: str):
        async with HttpAsyncHook(
            http_conn_id="pokeapi",
            method="GET",
        ).session() as session:
            response = await session.run(endpoint=url)
            return await response.json()

    get_pokemon.iterate(
        url=list_pokemon(),
    )

pokemon_iteration()

Comparison

Pattern Task Instances Work Per Task
get_pokemon.expand(url=urls) 100 1 Pokémon
get_pokemon.iterate(url=urls) 1 100 Pokémon

This demonstrates how Task Iteration can significantly reduce TaskInstance creation overhead. Task Spreading (running one iteration over exactly N TaskInstances with .batch(size=N).iterate(), to be renamed .spread()) is split out into #73688.

Notable design points addressed since the initial draft

  • .iterate()'s dict-argument semantics now match .expand(): passing a dict value forwards (key, value) pairs to each sub-task instead of bare keys.
  • IterableOperator.task_type and .operator_name both forward to the wrapped operator (including @task-decorated callables with a custom_operator_name), so sub-tasks report the correct type in the UI/API instead of always showing MappedOperator/IterableOperator.
  • XComIterable.flatten() moved out to Add XComIterable.flatten() to read an iterated task's pages as one sequence #73807, stacked on this PR, so this PR stays about running a task over its input.
  • on_kill() propagates to in-flight sub-tasks, and outlet/asset events recorded by a sub-task that already succeeded are replayed from its checkpoint on a later retry instead of being lost.
  • Deferred operators, reschedule-mode sensors, TriggerDagRunOperator, and ShortCircuitOperator-style downstream skipping are explicitly rejected inside IterableOperator with actionable errors rather than being silently mishandled — see the class docstring for the full list of current limitations.
  • multiple_outputs is ignored for iterated tasks, explicitly at the IterableOperator level. A @task with a Mapping return annotation infers multiple_outputs=True, but the value the runner pushes for an iterated task is the XComIterable aggregate rather than a dict, so honouring the flag made the runner reject the result after every sub-task had already succeeded. Each sub-task's return value is pushed whole as return_value_<index>; keys are not fanned out into separate XComs the way .expand() does. Documented on the class and in the Task SDK docs, and pinned by a runner-level regression test for a dict-returning task under .iterate().

Per-iteration keys: XComs and task state

Every iteration of an iterated task runs under the same task instance (same dag id, task id, run id and map index). Anything an iteration writes into a per-task-instance store therefore competes with its siblings for the same key, and with the async executor the winner is whichever iteration finishes last. Two stores are affected, and both now apply the same rule: a key written from inside an iteration carries that iteration's index.

  • XComs. IndexedTaskInstance.xcom_push/axcom_push suffix the key with _<index>. That is what makes return_value_<index> and XComIterable work, and it applies to any key an operator pushes from execute, including the keys of a multiple_outputs dict. Pulls are not suffixed: ti.xcom_pull(task_ids="upstream") reaches the upstream's XCom untouched.
  • Task state store (the AIP-103 store, context["task_state_store"]). This was a gap: an iteration that stored a watermark or a cursor with task_state_store.set("last_offset", ...) shared that key with every sibling. IndexedTaskStateStoreAccessor closes it: IndexedTaskInstance.task_state_store is the parent's accessor seen through the index, suffixing keys on get/set/delete and their async twins, and the sub-task's context carries the same object, so an operator does not need to know it is being iterated. clear() is refused inside an iteration, since it would wipe the siblings' state and the operator's own checkpoints; an iteration deletes its own keys instead.
  • The operator's checkpoints are separate. IterableOperator records per-index progress in the parent's store under _iterable_<index> and _iterable_completed, written through the parent's accessor, so they are never double-suffixed and never collide with user keys.

IndexedTaskRunner (formerly TaskExecutor, renamed because it read like one of Airflow's executors) builds the context an iteration runs against: a copy of the parent's context with the iteration's own task instance, its indexed state store view and its own outlet events. The operator binds it from the with statement, runs the operator inside that block, and records the outcome (checkpoint, XCom push, outlet-event merge) after it, so on_kill and the failure callbacks apply to the operator's execution only.


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

@dabla
dabla requested review from amoghrajesh, ashb and kaxil as code owners March 5, 2026 09:28
@dabla
dabla marked this pull request as draft March 5, 2026 09:35
@dabla
dabla force-pushed the feature/dynamic-task-iteration branch 3 times, most recently from d8a30b9 to edad5de Compare March 5, 2026 12:39
kaxil
kaxil previously requested changes Mar 5, 2026

@kaxil kaxil left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for working on this — excited to see DTI taking shape for Airflow 3.2. I've gone through the full diff and have feedback on the implementation, some are bugs that would crash at runtime, others are design choices worth iterating on.

A few high-level things:

  1. No tests. ~700 lines of new production code with zero test coverage. We need tests for IterableOperator, TaskExecutor, MappedTaskInstance, HybridExecutor, XComIterable, DecoratedDeferredAsyncOperator, and the iterate/iterate_kwargs methods — covering success, failure, retry, deferral, and edge cases.

  2. Worker resilience. Since DTI runs N sub-tasks inside a single worker process, we need to think through what happens when that worker dies mid-execution — the scheduler has no record of which sub-tasks completed. Worth documenting the expected behavior and trade-offs here (and whether we want to add checkpointing later).

  3. Thread safety. Several shared mutable structures (context dict, os.environ) are accessed concurrently from multiple threads without synchronization. This needs to be addressed before merge.

Inline comments below with specifics.

Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py Outdated
Comment thread task-sdk/src/airflow/sdk/bases/operator.py Outdated
Comment thread task-sdk/src/airflow/sdk/definitions/mappedoperator.py
Comment thread task-sdk/src/airflow/sdk/definitions/_internal/expandinput.py
Comment thread task-sdk/src/airflow/sdk/definitions/_internal/expandinput.py
Comment thread task-sdk/src/airflow/sdk/execution_time/executor.py
Comment thread task-sdk/src/airflow/sdk/execution_time/lazy_sequence.py Outdated
Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py Outdated
Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py Outdated
Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py Outdated
@dabla

dabla commented Mar 5, 2026

Copy link
Copy Markdown
Contributor Author

Thanks for working on this — DTI is an interesting concept and I can see the use case. I've gone through the full diff and have a number of concerns, some are bugs that would crash at runtime, others are architectural questions worth discussing before this goes further.

A few high-level things:

  1. No tests. ~700 lines of new production code with zero test coverage. We need tests for IterableOperator, TaskExecutor, MappedTaskInstance, HybridExecutor, XComIterable, DecoratedDeferredAsyncOperator, and the iterate/iterate_kwargs methods — covering success, failure, retry, deferral, and edge cases.

Thanks for pointing this out. As mentioned earlier on Slack, this PR is currently intended as an initial draft to demonstrate the concept and gather early architectural feedback.

I agree that proper test coverage is essential before this can move forward. The plan is to add unit tests covering the components you mentioned (IterableOperator, TaskExecutor, MappedTaskInstance, HybridExecutor, XComIterable, DecoratedDeferredAsyncOperator, and the iterate/iterate_kwargs APIs), including scenarios for success, retries, failures, deferral, and edge cases.

Once we converge on the architectural direction, I will add the corresponding test suite.

  1. Architectural concern. This builds a mini-executor inside an operator — running N tasks in threads with in-memory XCom, custom retry logic, and sleep()-based retry delays. The scheduler has no visibility into sub-task states, so if the worker dies mid-execution there's no record of which sub-tasks completed. This feels like it needs broader design discussion (probably an AIP) before merging, since it fundamentally changes how task execution works.

I agree this is an important architectural concern and worth discussing further.

The goal of this prototype is to explore a trade-off between observability and scheduling overhead, @ashb and @potiuk mentioned the same remark before. If we try to preserve the same visibility and lifecycle guarantees as Dynamic Task Mapping, we essentially end up re-implementing DTM semantics, which brings back the same scheduler overhead that this approach is trying to avoid.

This proposal intentionally explores a different point in that trade-off space: executing iterations within a single task while allowing controlled parallelism. That does mean the scheduler has indeed less visibility (but also less load) into the internal execution units.

  1. Thread safety. Several shared mutable structures (context dict, os.environ) are accessed concurrently from multiple threads without synchronization.

Good point — thread safety needs to be handled carefully here.

Regarding the task context, my understanding is that operators already receive a per-task context instance, but you're right that when running iterations concurrently we should avoid sharing mutable structures across threads. One possible approach would be to create a shallow or deep copy of the context for each execution unit to ensure isolation.

If you have concerns about specific structures (e.g., os.environ or others), I'm happy to address them and introduce appropriate synchronization or isolation mechanisms where needed.

@dabla dabla changed the title refactor: Implemented Dynamic Task Iteration Implemented Dynamic Task Iteration Mar 5, 2026
@kaxil
kaxil self-requested a review March 12, 2026 00:02
@kaxil
kaxil dismissed their stale review March 12, 2026 00:02

Stale review

Comment thread task-sdk/src/airflow/sdk/bases/operator.py Outdated
Comment thread task-sdk/src/airflow/sdk/bases/operator.py Outdated
Comment thread task-sdk/src/airflow/sdk/bases/operator.py Outdated
Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py Outdated
Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py
Comment thread task-sdk/src/airflow/sdk/execution_time/executor.py
Comment thread task-sdk/src/airflow/sdk/definitions/_internal/expandinput.py
Comment thread task-sdk/tests/task_sdk/definitions/conftest.py Outdated
@dabla
dabla force-pushed the feature/dynamic-task-iteration branch from 960438c to 765fcfb Compare March 18, 2026 23:10
@dabla
dabla marked this pull request as ready for review March 19, 2026 17:05
@dabla
dabla requested a review from kaxil March 19, 2026 20:27
@kaxil
kaxil requested a review from uranusjr March 20, 2026 00:19
Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py Outdated
Comment thread task-sdk/src/airflow/sdk/bases/operator.py Outdated
Comment thread task-sdk/src/airflow/sdk/definitions/mappedoperator.py Outdated
Comment thread task-sdk/tests/task_sdk/definitions/conftest.py Outdated
@kaxil

kaxil commented Mar 20, 2026

Copy link
Copy Markdown
Member

@uranusjr You should also review this PR since it touches several important modules :)

@dabla
dabla force-pushed the feature/dynamic-task-iteration branch from b11f852 to 9f2c750 Compare March 20, 2026 08:35
@dabla
dabla requested a review from kaxil March 20, 2026 16:01
@dabla
dabla force-pushed the feature/dynamic-task-iteration branch from 16ec1fc to 3242037 Compare March 20, 2026 17:59
@dabla dabla changed the title Implemented Dynamic Task Iteration AIP-98: Dynamic Task Iteration Mar 21, 2026
@dabla

dabla commented Mar 21, 2026

Copy link
Copy Markdown
Contributor Author

@uranusjr @kaxil In our patched Airflow installation I had to register the XComIterable manually with serde for serialization. How do I make sure it’s automatically registered with serde?

@dabla
dabla marked this pull request as draft April 14, 2026 19:43
@dabla dabla changed the title AIP-98: Dynamic Task Iteration AIP-104: Dynamic Task Iteration and Dynamic Task Partitioning Apr 17, 2026
@dabla

dabla commented Oct 8, 2026

Copy link
Copy Markdown
Contributor Author

@dabla To increase the odds of getting this into 3.4 my request to you is to thoroughly review the code before pushing and testing e2e.. otherwise the back-and-forth consumes too much time

I always run my test dags with iterable and task spreading locally in my devcontainer to make sure it still works, the code is, in the best I can, thoroughly reviewed and apply the design patterns there where the add value.

@kaxil

kaxil commented Oct 8, 2026

Copy link
Copy Markdown
Member

Kaxil's case:

#62922 (comment) :) This feels like a trivial one that can be avoided easily

@dabla

dabla commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor Author

Kaxil's case:

#62922 (comment) :) This feels like a trivial one that can be avoided easily

Yes added that as a rule in my skill to make sure it doesn't happen anymore ;-)

dabla and others added 27 commits October 9, 2026 12:43
Add Iterable Tasks: `.iterate()` and `.iterate_kwargs()` on operators and
`@task`, the counterpart of `.expand()` that processes every item inside
one task instance instead of creating one task instance per item.

- IterableOperator resolves the input by index as `.expand()` does and
  runs the items on AsyncAwareExecutor: sync operators in a thread pool,
  async operators concurrently on one event loop, up to
  `task_concurrency` at a time.
- Each item's return value is pushed as `return_value_<index>`, and the
  task returns an XComIterable, a lazy read-only Sequence over them that
  a downstream `.expand()` or `.iterate()` consumes. Skipped items are
  left out, and downstream tasks with `all_success` are skipped, as with
  a mapped upstream.
- Per-item progress is checkpointed in the task state store (AIP-103),
  tied to the item's input and the attempt that wrote it, so a retry or a
  clear after a failure resumes the items that already succeeded and a
  clear after success runs them all again. Outlet events are replayed
  from the checkpoint.
- XComs and task state written from an item carry its index, and each
  item runs against its own view of the context.
- Deferral, reschedule-mode sensors, TriggerDagRunOperator and
  downstream skipping from an item are rejected with a clear error.
- Documented in task-sdk/docs/mapped-tasks-vs-iterable-tasks.rst.

Co-Authored-By: Tzu-ping Chung <uranusjr@gmail.com>
Co-Authored-By: Copilot <223556219+Copilot@users.noreply.github.com>
Two tests of the iterable operator named the reviewer whose thread led to
them instead of saying what they check. A test's name and docstring must
stand on their own: one now says that a succeeded item's outlet events are
merged into the task's accessor, the other that a failed item's retry
callback waits until no sibling rules the retry out.
The item's view of the context swapped ti, task_instance, task_state_store
and outlet_events but not task, so inside execute and in the item's
callbacks context["task"] was the IterableOperator. Under .expand() it is
the unmapped operator of the task instance: context_update_for_unmapped
sets context["task"] next to ti.task, and the rendering in _create_task
already got that. IndexedTaskRunner.indexed_context now swaps the same key.

Two tests fail before the change: the runner test for the indexed context
asserts the item's operator under "task" and the parent's own left alone,
and an operator test checks that execute and on_success_callback of every
item see the unmapped operator there, the same object as ti.task.
The docs page presented iteration as a way to keep large results out of
the metadata database through a custom XCom backend. Each item's result is
also written to its checkpoint, which goes to the task_state_store table
for the store's retention unless a state_store_backend is configured, so
with only a custom XCom backend the payloads move from one table to the
other. The triggerer bullet, the "XCom backend" row of the comparison
table and the class docstring now say so and point at state_store_backend.
A task instance sends one event per asset, so items of an iterated task
that emit to the same asset are merged into that one event: extra keeps
what the last item to finish wrote, partition keys and alias events
accumulate. .expand() sends one event per mapped task instance, and a
consumer reading triggering_asset_events sees all of them. The comparison
table gets an "Asset events" row and _merge_outlet_events a note, since a
per-file Metadata pattern ported from .expand() produces one event here.
.expand() never sees a string or another scalar from upstream: the
upstream's _push_xcom_if_needed raises UnmappableXComTypePushed when it
has a mapped dependant. An IterableOperator is not a MappedOperator, so
iter_mapped_dependants does not find it and that check never fires; an
upstream returning a JSON string made .iterate() succeed over one wrong
item. Source.from_argument now applies the same rule to the value it
resolved from an XComArg, with the is_mappable_value the push check uses,
and XComForMappingNotPushed for None, as the push would have raised.

Tests cover strings, bytes, scalars and None on Source, and .iterate()
over an upstream returning a JSON string, which fails before any item runs.
With K failed items and a retry_policy, the policy was evaluated K times
in _failure_for_the_runner to choose the exception for the runner, once
more in _task_will_retry to pick the failed items' callbacks, and once
more by the runner itself. A policy that calls a model, such as
common.ai's LLMRetryPolicy, pays for each call and may answer differently
each time, so the callbacks could announce a retry the runner then did not
take, which deferring them was meant to prevent.

_failure_for_the_runner now keeps the decision taken for the exception it
chose, and _task_will_retry reuses it; the policy is only evaluated there
for an exception that was never chosen (a single failure). Two tests fail
before: a counting policy is evaluated once per failed item, and a policy
that alternates its answers gives the callbacks the same outcome as the
exception handed over.
…ff the loop

Only run() went to the thread pool for a sync item: IndexedTaskRunner's
enter and exit, and so the item's on_success_callback and
on_skipped_callback fired from the exit, ran on the event loop thread. A
notifier reading a connection or a variable there while a sibling was
inside an async SDK call got DeadlockImminentError, a BaseException that
_run_task_state_change_callbacks does not catch: the item had done its
work, was checkpointed UP_FOR_RETRY anyway, and the task failed without a
retry with a message blaming an async sub-task. The user wrote an ordinary
sync operator; the operator put it on a thread.

A sync item now enters and exits its runner inside the function handed to
the thread pool, so its callbacks run where a sync SDK call waits for the
comms lock; async items stay on the loop. A sync item whose coroutine was
cancelled while its thread went on reports nothing from its exit
(IndexedTaskRunner.cancel), since it gets no checkpoint and runs again on
the next attempt. on_kill catches BaseException per sub-operator, so one
failure no longer skips the rest, and runs the kills off the loop thread:
in a thread of their own when the SIGTERM handler calls it while the loop
runs, through asyncio.to_thread from _run_tasks. The failure message for a
DeadlockImminentError says which thread the call was made on, and the
class docstring says where each callback runs.

Tests pin the thread of execute and of every callback for a sync and an
async item, the kill reaching the second sub-operator when the first
raises DeadlockImminentError, the kill running off a running loop, and a
cancelled sync item firing no callback when its thread finishes.
On SIGTERM the runner calls on_kill() once and the items in flight are
killed, but nothing stopped the iteration from pulling more: each killed
item came back as a failure and freed its slot, _fill_pending submitted
the next item into it, and new items with their remote jobs kept starting
until the supervisor escalated to SIGKILL. Those never got an on_kill.

IterableOperator keeps a stop flag that on_kill() sets first.
AsyncAwareExecutor.imap_unordered takes it as ``stop`` and asks it before
and after every pull: once set nothing more is submitted, an item pulled
at that moment included, and what was already submitted drains. After the
loop, before the other outcomes, _run_tasks raises AirflowTaskTerminated
saying how many items ran and that the rest never started: killed items
returning normally would otherwise have written the completion marker and
returned an XComIterable over XComs that do not exist. The task fails
without a retry, as the runner treats a terminated task, and no marker is
written, so a later clear resumes from the checkpoints of the items that
did finish.

Tests: the executor stops pulling once ``stop`` says so and still drains
what it submitted; a kill from inside the first of four items runs no
other item, leaves no marker and only that item's checkpoint; a kill
while the input is resolved runs no item and is not an empty-input skip.
IterableOperator carried seven attributes for what one run remembers
while it is going: the sub-operators in flight with their lock and the
ones already killed, the stop flag, the runners of the failed items, the
resolved input and the retry policy's decision for the exception handed to
the runner, plus a __deepcopy__ override to give every copy fresh ones.

They move into IterationState, in the same module. A copy of the state is
a fresh one, so the override goes, and the operator starts every run with
a fresh state too. The state exposes what the operator means rather than
its containers: register/unregister and take_in_flight for the kill,
request_stop/stop_requested for the executor's stop, note_failed and the
failed_runners tuple, keep_decision/decision_for, and length for the
resolved input; the operator makes no reach-through call into it.
IndexedTaskRunner registers through a SubOperatorRegister protocol instead
of a dict and a lock.

Tests cover the state on its own: registration until unregistered, keying
by identity for operators that compare equal, take_in_flight handing each
operator out once, the stop flag, the order of failed runners, a kept
decision answering for that exception object only, length unknown until
resolved, and a deep copy being fresh while the original keeps its state.
BaseOperator.__init__ runs under _apply_defaults, which fills every
parameter of its signature the call leaves out from the DAG's default_args.
IterableOperator forwarded the retry, scheduling and pool settings of the
wrapped operator but left the five on_*_callback parameters and
pre_execute/post_execute out, so a DAG with
default_args={"on_failure_callback": notify} gave the iterated task that
callback too, while the items got it through the wrapped operator's partial
kwargs: with three failed items notify ran four times, the last one with
the parent's context, where .expand() and .partial(on_failure_callback=
notify) run it three times. The class docstring and the docs page say the
iterated task carries no callbacks of its own.

They are now passed as None explicitly, so _apply_defaults leaves them
alone and both paths agree with the docs. The test builds a DAG with all
seven in default_args and checks the iterated task has none while an
unmapped item has them; it fails before the change.
on_kill() took a thread only when the event loop was running on the calling
thread. The runner's SIGTERM handler can also arrive while the loop is
paused between two run_until_complete calls, as imap_unordered hands a
result to the consumer; the kills then ran inline on the main thread, and a
synchronous SDK call in a sub-operator's on_kill waited for the comms lock
a parked asend held, which only the paused loop could release: the loop
never resumed, the remaining on_kills never ran, and SIGKILL ended it.

on_kill() now always sets the stop flag, takes the in-flight snapshot and
kills in a daemon thread of its own, where the call waits its turn in both
cases while the loop goes on. _run_tasks no longer goes through on_kill():
it kills through asyncio.to_thread as before, awaited while the loop runs.

A new test calls on_kill() with no loop running on the main thread, as in
the paused window, and checks the kill ran off it; the tests that asserted
the kill right after the call wait for the thread.
…nd the async item's own timeout too

Two on_kill call sites still ran on the loop thread. IndexedTaskRunner.in_flight
killed the operator the parent's execution timeout struck as the timeout unwound,
and _execute_async_task killed an async item that ran out of its own
execution_timeout; both on the loop thread, where an on_kill that cancels a
remote job through a sync hook raises DeadlockImminentError. The first site let
it through its `except Exception`, the second was unguarded, so the parent's
timeout became an AirflowFailException: no retry, and the job left running.

The operator the parent's timeout strikes now stays registered as the timeout
unwinds (the timeout only lands on async items: the signal reaches the main
thread, where the loop runs, while sync items sit in worker threads), and
IterableOperator._run_tasks kills it off the loop with the others through
asyncio.to_thread. An async item's own timeout kills it through
`await to_thread(task.on_kill)`.

Tests fail before: the runner leaves the struck operator registered and does
not kill it in place; the async item's own timeout kills it off the loop thread;
an async item raising the parent's timeout inside its coroutine, with an
on_kill that raises on the loop thread, still ends in AirflowTaskTimeout with
both items killed.
…ync item

Every sync item with an execution_timeout went through _run_execute_callable
in its worker thread, which sent SetExecutionTimeout to the supervisor again
and tried TimeoutPosix, which only works on the main thread: the supervisor's
hard-kill deadline drifted to "timeout after the last item started", up to a
full timeout late, and the TimeoutPosix warning was logged once per item.
With 4 sync items and a 30 s timeout: 4 sends, 4 warnings.

_run_execute_callable takes enforce_timeout, False for an IndexedTaskInstance:
the items run in worker threads under the parent's own limit, which the
parent enforces and reported once. The async items already handled their
limit themselves. A test runs three sync items with a 30 s timeout through
the parent's _run_execute_callable and checks a single SetExecutionTimeout;
it sees four before the change.
The callbacks section of the docs page and the class docstring describe which
callbacks run per item; neither said what happens to listeners. They fire once,
for the task instance, when the runner reports its state, as for any task: an
item is not a task instance and fires none, where .expand() fires them once per
mapped task instance.
IndexedTaskInstance.xcom_push writes under <key>_<index>, so iterations never
overwrite each other, but a pull had no counterpart since the override that
suffixed every pull was removed: ti.xcom_pull(key="progress") inside an item
read the parent's unsuffixed key and got None for the value the item had just
pushed, and reading it back needed key=f"progress_{ti.index}". The task state
store already follows the symmetric rule: its accessor suffixes get, set and
delete alike.

xcom_pull and axcom_pull now add the index for a pull of the iteration's own
XComs, when no task is named or its own task in its own DAG is, and leave a
pull from another task, from several, or from another DAG untouched, which is
what the review of 2026-09-10 asked for. The docs page and the class docstring
say so, and how to read another iteration's value.

Tests fail before: which key reaches the base method for six kinds of pull
and for the async pull, and an item that pushes a key and reads it back
through both forms of its own pull.
_run_tasks told apart, inline, the item outcomes the iteration cannot carry
(a deferral, a reschedule, a DAG run trigger, a downstream skip, a
BaseException that is no Exception) from the failures it collects, in a chain
of isinstance checks that grew with every round. They move into
IterableOperator._fail_fast_for, which returns the AirflowFailException that
ends the task at once for such an outcome, or None; the loop raises what it
returns and collects the rest. Same messages, same behaviour; the operator's
own policy stays on the operator rather than on the expand input, since the
exceptions come from running the items, not from resolving the input.
The exception handling of the iteration was spread over four protected
methods of IterableOperator and a 95-line _run_tasks, with the failed
runners and the kept retry policy decision living on IterationState.
IndexedTaskOutcomes now holds all of it: a context manager entered next
to Checkpoints around the loop, whose record() takes each indexed task's
outcome (count, skip, collect, or raise at once for an outcome the
iteration cannot carry) and whose conclude() raises the task's own
outcome from what was recorded and from the run's IterationState (killed,
failed on the exception the runner judges, empty input, every indexed
task skipped). Whatever exception leaves the block, the failed indexed
tasks' callbacks are reported on exit with the task's fate, so a
fail-fast raised from record() and the outcomes raised by conclude()
take the same path. IterationState keeps only what the kill needs.

The prose of both modules says "indexed task" where it said "item", and
_run_task takes the run's collector as a keyword argument.
The module-level functions of iterableoperator.py each served one class.
The constructor checks and the input fingerprint pair are now methods of
IterableOperator; the outlet-event snapshot and its replay belong to the
checkpoint, IndexedTaskState, which owns that field; the live merge of
one indexed task's events into the parent's belongs to IndexedTaskRunner,
which owns the accessors. The replay writes into the parent's accessors
directly instead of building a throwaway set and merging it. The retry
decision weights are an attribute of IndexedTaskOutcomes, their only user.
…ill that arrives before the run

IndexedTaskStateStoreAccessor wraps the parent's accessor without calling
its constructor, so the inherited _clear_backend_only read a _scope the view
never had; it now refuses as clear() does. The run's IterationState is no
longer replaced when _run_tasks starts, which discarded an on_kill() that
landed before the run: execute() renews it once the run ended, so a kill
before the run stops it and a rerun in the same process starts clean.
XComIterable.deserialize returns cls(**data); a conftest comment reads again.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants