Skip to content

AIP-104: Task Iteration - #62922

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

dabla wants to merge 10 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 Outdated
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 mentioned this pull request Sep 24, 2026
@dabla

dabla commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor Author

Scope change: this PR is now Task Iteration only.

Following the feedback from today's Airflow dev call, I have split AIP-104 in two so this PR stays focused on .iterate() / .iterate_kwargs():

  • AIP-104: Task Iteration #62922 (this PR), "AIP-104: Task Iteration": IterableOperator, AsyncAwareExecutor, XComIterable, the async resolution of expand inputs and XComArgs, the runner's sub-task execution with checkpoints, outlet-event replay and on_kill propagation. Head is a3d7f65.
  • AIP-104: Task Spreading #73688, "AIP-104: Task Spreading" (draft): .batch(size=N).iterate(...) (to be renamed .spread()), MappedIterableOperator, BatchedExpandInput, the runtime size resolved by the scheduler from mapped_length, and the related docs and tests. It is stacked on this PR, so its diff includes this one until this merges; the spreading work is its last commit. The naming question and the devlist vote move there too.

The split commit is a3d7f65. Besides removing spreading it changes one thing in what stays here: .iterate() no longer funnels through a size=0 BatchedOperator. OperatorPartial and _TaskDecorator build the IterableOperator directly, and _expand owns the mapped-operator construction on both classes again, returning the operator so the public methods wrap it. The batching PR builds on those helpers instead of undoing them. Everything discussed in the resolved threads on this PR is unchanged apart from that.

The description and title are updated accordingly. Suggested reading order for the remaining diff: iterableoperator.py, then task_runner.py, then executor.py, then xcom.py, then the async plumbing in expandinput.py, xcom_arg.py and lazy_sequence.py.


Drafted-by: Claude Fable 5.1; reviewed by @dabla before posting

@dabla dabla mentioned this pull request Sep 24, 2026
@dabla

dabla commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

Rework: .iterate() now resolves its input by index, the way .expand() does

Commits 79e7a33 and 9a75b78. The goal is to make this PR easier to review, to simplify the code, to avoid unnecessary duplication, and to reuse what already exists as much as possible.

Why. IterableOperator walked its input through iter_values()/aiter_values(), a streaming cross product that existed nowhere else in Airflow and had to be kept in a sync and an async flavour so the two could not drift. The length was only known after the stream had been drained, recorded as a side effect by count()/async_count(), and the sync flavour was never called at run time. Every input .iterate() accepts is a sized sequence once pulled, so it can do what .expand() already does: know the length up front and pick each item by index.

What changed.

  • ExpandInput.aresolve(context) pulls every source once and returns Resolved(length, aget). DictOfLists is the cross product, ListOfDicts the mapping at that position, Decorated wraps it in op_kwargs.
  • The cross-product arithmetic is extracted into index_for_each_field(), now shared by _expand_mapped_field() (.expand()) and aresolve() (.iterate()) and tested on its own, so both hand a sub-task the same item for a position.
  • Source reads one argument by index without blocking the loop thread: a value's own alen()/aget() when it has them (XComIterable, LazyXComSequence, which gains an alen() twin of __len__), else a worker thread.
  • XComIterable carries a flattened_length that IterableOperator.axcom_push tallies from each value as it is pushed, so FlattenedXComIterable knows its length without reading a page. It is now a cursor over the last page fetched: a sequential read fetches every page exactly once, where indexing the old class re-walked the pages from the start on every call.
  • Removed from the PR's API surface: XComArg.iter_values/aiter_values, ResolveMixin.iter_values, the ExpandInput sync/async twins, count(), async_count(), aiterate() and _to_iterable().

Net effect: 150 fewer lines of source.

Behaviour change. iterate_kwargs([xcom1, xcom2]) now follows expand_kwargs, one mapping per XComArg, where it used to flatten each into several.

Supervisor calls. Unchanged for in-memory and plain XCom inputs. A mapped upstream costs one GetXComCount per source, since the length is needed up front; its per-item reads are unchanged. A cross product no longer re-resolves a later source once per item of an earlier one.

Note on the static checks job. The failing hook regenerates the Java SDK dependency verification metadata; this PR does not touch the Java SDK. The branch is behind main, whose newer Java build the metadata file already reflects, so the hook prunes entries the branch's build does not reference. It clears on the next rebase onto main.

@dabla

dabla commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

Follow-up for the one cost this rework leaves on the table: a mapped upstream (LazyXComSequence) is still read one item per request. #73790 makes those reads chunked (one GetXComSequenceSlice per 32 items, on both the sync and the async path) and is stacked on this PR, so the per-item reads stay as they are here and this PR remains a pure simplification.

@dabla

dabla commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

Brought level with main: main is merged in just below the rework, which now sits on top as fe3e7bc and e7f39eb. A commit-by-commit rebase of the branch's history conflicted from its first commit against work main already contains, so the merge keeps the reviewed history intact with the same resulting tree. This also clears the Java SDK verification-metadata hook failure in the static checks.

@dabla

dabla commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

XComIterable.flatten() moved out to its own stacked PR, #73807, with its docs (dbf2bdc here). Task iteration does not need it: nothing in the SDK outside bases/xcom.py referred to it, so this PR is now only about running a task over its input. About 320 lines lighter. The stack is now #62922 → #73688 (spreading), #73790 (chunked lazy reads), #73807 (flatten), each one commit on top of this branch.

@dabla

dabla commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor Author

#62922 — AIP-104 Task Iteration (base of the stack)
.iterate() runs a task over its input inside one task instance, with an async executor. Reworked so it resolves its input by index the way .expand() does, with the cross-product rule in one shared helper. No longer adds iter_values/aiter_values to the SDK. Level with main. Everything below is one commit on top of it and rebases once it merges.

#73688 — AIP-104 Task Spreading
.batch(size=N).iterate(...): the same iteration spread over N task instances, items dealt round-robin. The scheduler fixes N before the task runs; each instance takes its share with a range over the resolved input. Naming still open, .spread() leads the devlist vote.

#73790 — Chunked reads of a mapped upstream's XCom sequence
LazyXComSequence fetched one item per request. Now one GetXComSequenceSlice per 32 items on both the sync and async paths, at most one chunk held, negative indices through the cached count. For a 17,000-item mapped upstream consumed by .iterate(), about 530 requests instead of 17,000. Chunk size is a constant for now.

#73807 — XComIterable.flatten()
An iterated task whose iterations return pages of items can hand them downstream as one sequence. The producer tallies the flattened length as it pushes, so the view knows its length without a read, and it holds one page at a time. Split out of #62922 because iteration does not need it, and documented here.

@dabla

dabla commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

#70223 (the POST …/xcoms/{dag_id}/{run_id}/{task_id}/keys endpoint that lets XComIterable iterate and slice in one round trip) is now stacked on this PR as well, squashed to one commit, f688623. The stack is: #62922 → #73688 (spreading), #73790 (chunked lazy reads), #73807 (flatten), #70223 (batched XComIterable reads), each one commit on top of this branch with a review-only link at the top of its description.

@dabla

dabla commented Sep 28, 2026

Copy link
Copy Markdown
Contributor Author

Three more commits, following up on the earlier remark about suffixing on the pull side and a gap it uncovered on the state store:

  • f5deadd Give each iteration its own task state store keys, as for XComs. Every iteration runs under the same task instance, so a key one iteration stored through context["task_state_store"] was shared with its siblings, last writer wins. IndexedTaskStateStoreAccessor decorates the parent's accessor and suffixes keys with the index on get/set/delete, sync and async, the same rule IndexedTaskInstance.xcom_push applies to XComs. clear() is refused inside an iteration. The operator's own checkpoints keep their unsuffixed _iterable_<index> keys in the parent's store. Tested on the accessor, on IndexedTaskInstance, and end to end with a sync and an async operator that keep state.
  • 55c48ed Build the sub-task's context in the runner. The clone of the parent's context with the sub-task's ti, state store view and outlet events was built twice in the operator; it now lives in one place next to the task instance it belongs to, and the two _run_operator methods are gone.
  • 7dd2acc Rename TaskExecutor to IndexedTaskRunner. It read like one of Airflow's executors; it runs one index of an iterated task, so it is named after IndexedTaskInstance.

The four stacked PRs are restacked on top; each still one commit.

@dabla

dabla commented Sep 28, 2026

Copy link
Copy Markdown
Contributor Author

Small amend: the rename commit is now 742d8b1, binding the runner straight from the with statement as indexed_task_runner. The stacked PRs are restacked accordingly.

@dabla

dabla commented Sep 28, 2026

Copy link
Copy Markdown
Contributor Author

Level with main again (merge efd32cb below the PR commits, which are re-applied on top; tip 46572c7). The four stacked PRs are restacked, one commit each.

@dabla

dabla commented Sep 28, 2026

Copy link
Copy Markdown
Contributor Author

Level with main again (16 commits merged below the PR commits; tip 85473e5). The four stacked PRs are restacked, one commit each.

@dabla

dabla commented Sep 29, 2026

Copy link
Copy Markdown
Contributor Author

6579f20 moves IndexedTaskRunner from executor.py to task_runner.py, next to the IndexedTaskInstance it runs and the _execute_task/_execute_async_task it calls, so executor.py is only AsyncAwareExecutor again. indexed_context() is now a context manager: it builds the indexed task's view of the context, remembers it for the state-change callbacks and makes it the current context for the block. Its tests moved with it. Stacked PRs restacked.

@dabla

dabla commented Sep 29, 2026

Copy link
Copy Markdown
Contributor Author

Level with main again (13 commits merged below the PR commits; tip 21b1794). The four stacked PRs are restacked, one commit each.

Comment thread task-sdk/src/airflow/sdk/execution_time/task_runner.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
Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py

.. _sdk-dynamic-task-mapping-vs-iteration:

Dynamic Task Mapping vs Iterable Tasks

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.

Ash mentioned on Sept 17, in the Slack thread about renaming the AIP, that he wants to softly rename "Dynamic Task Mapping" to just "task mapping" or "mapped tasks", since Loops and the dynamic execution graph are more dynamic than mapping is. This page still uses "Dynamic Task Mapping (DTM)" throughout: the title here, the headings at lines 67 and 370, the comparison table, and "DTM" as shorthand across the prose. Since the page is new, this is the cheap moment to follow that, including the file name and the sdk-dynamic-task-mapping-vs-iteration label that deferred-vs-async-operators.rst links to. Something like "Mapped tasks vs iterable tasks" would read well.

list_pokemon_task >> get_pokemon_task


The scheduler only manages a single task. With sync tasks, iterations are

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.

This paragraph has it backwards. Threads overlap blocking I/O fine up to task_concurrency: the sync HttpOperator example above, with the request stubbed to sleep, ran 8 requests at once and finished about seven times faster than serially. Pure-Python CPU-bound work gets no speedup from threads under the GIL, which also contradicts "avoid Iterable Tasks when each item represents a long-running or heavy computation" further down. Maybe: threads overlap blocking I/O up to task_concurrency, async scales further because coroutines are cheaper than threads, and CPU-bound work doesn't speed up. A few smaller claims on the page are off too. The benchmark table at the top compares mapped SFTPOperator with hand-written @task loops, none of which use .iterate(), so it should be labelled that way or get an .iterate() row. Line 225 says "For 5 Pokémon" but the example fetches limit=100. Line 258 says iterations share the event loop "(and connection)", but nothing shares a connection, since HttpAsyncHook.session() opens a new session per call.

Comment thread task-sdk/src/airflow/sdk/definitions/iterableoperator.py Outdated
them, so a ``dict`` return annotation on the task does not fan its keys out into separate
XComs the way it does with ``expand()``. Every key an iteration pushes or stores carries its
index the same way: ``ti.xcom_push("foo", v)`` in iteration 2 lands under ``foo_2``, and so
does ``task_state_store.set("foo", v)``, so iterations never overwrite each other's values.

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.

This holds inside the iteration's own thread or coroutine, but a threading.Thread the task starts, or loop.run_in_executor, doesn't carry the context over. get_current_context() there falls back to the parent's context with its un-indexed ti and task_state_store, so keys pushed from such a thread get no suffix and iterations overwrite each other: with iterate(x=[1, 2, 3]) and a push from a helper thread, only one row survived. A sentence here pointing at asyncio.to_thread or contextvars.copy_context().run(...), both of which do carry it, would cover it.

entered.append(index)
if len(entered) == 2:
both_entered.set()
await both_entered.wait()

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.

The second task finds both_entered already set, so its wait() returns without yielding, and it reads its context and leaves its block before iteration 0 resumes. The reads happen in nesting order, so swapping the ContextVar for a plain thread-local stack still passes this test (and the rest of the context, runner and iterate suites). An await asyncio.sleep(0) after the wait, with both tasks reading while the other is still inside its block, would make it catch that regression.

@kaxil

kaxil commented Sep 29, 2026

Copy link
Copy Markdown
Member

Can you squash commits please?

dabla and others added 9 commits September 30, 2026 07:29
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/dynamic-task-mapping-vs-iteration.rst.

Co-Authored-By: Tzu-ping Chung <uranusjr@gmail.com>
Co-Authored-By: Copilot <223556219+Copilot@users.noreply.github.com>
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The error an iterated async task gets for a synchronous SDK call already
points at Variable.aget/aset, which apache#72329 adds; the class docstring
still said Variable had no async equivalent.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
partial() already prefixes the wrapped operator's task id with the task
group, and BaseOperator.__init__ prefixed it again for the
IterableOperator, so inside a TaskGroup the task was registered as
"tg.tg.f" while its items pushed their XComs for "tg.f", a task id with
no task instance. The IterableOperator now gets the bare id, and an
iteration takes the task id of the task instance that runs it.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
BaseOperator.__deepcopy__ calls copy.copy on every attribute in
shallow_copy_attrs, which held the lock guarding the sub-tasks in
flight, and a lock cannot be copied: deepcopy of an iterated task and
dag.partial_subset() failed on any Dag using .iterate(). A copy is
another task with nothing in flight, so it now gets a fresh lock and an
empty set through the deepcopy memo.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The execution timeout unwinds through the executor, which cancels every
item's coroutine, and each one left the register of operators in flight
before the runner called on_kill(); sync items too, since they were
registered in the coroutine rather than in the thread still running
execute. The register was also a set, and the sub-operators of one
iterated task compare equal, so it never held more than one of them.
on_kill() now runs before the executor cancels, operators are
registered where their code runs and keyed by identity, the item the
timeout strikes directly is killed as it unwinds, and each is killed
once although the runner calls on_kill() again.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
create_indexed_task built the indexed task instance without the
parent's _ti_context_from_server, so every iteration had no logical
date, a template context without dag_run or ds, and get_previous_ti()
and get_previous_dagrun() answered for no run at all.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Every item failure reached the runner inside a BaseExceptionGroup, and
the runner decides by exception type: a retry_policy rule never matched
the group, and AirflowSensorTimeout, which the runner fails without a
retry, was retried. Fail-fast exceptions are now raised on their own, a
single failure unwrapped, and with several failures the retry policy is
evaluated on each item's exception and the one whose decision weighs
most is raised, the others attached as its cause. IndexedTaskRunner
treats the fail-fast exceptions as final for the callbacks too.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
An iteration is unmapped from the MappedOperator, whose downstream task
ids are empty because the edges land on the IterableOperator, so a
ShortCircuitOperator took its "no downstream tasks" early return and a
branch operator found nothing to skip: every downstream task ran.
.iterate() now refuses any SkipMixin operator when the Dag is defined.
The check is on the class, since the @task path's
can_skip_downstream is False even for @task.short_circuit and
@task.branch.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The SUCCESS checkpoint was written before the result was pushed to XCom,
and a failing push fell through to the handler that overwrites it with
UP_FOR_RETRY, so the retry ran the operator again for work that had
finished. Publishing now fails on its own: the checkpoint stays, the
task retries, and the retry replays the result from the checkpoint.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@dabla

dabla commented Sep 30, 2026

Copy link
Copy Markdown
Contributor Author

Can you squash commits please?

Squashed commit like you asked, all current commits are related to each comment of the last review round.

The runner deletes every XCom of the task before an attempt, and a
retry that skips an item which already succeeded replayed only its
return value and outlet events, so any other key it pushed was lost
although the task then succeeded. The keys an item pushes are now kept
in memory while it runs, written once with its SUCCESS checkpoint and
pushed again when a retry skips it.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

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