dabla commented on code in PR #62922: URL: https://github.com/apache/airflow/pull/62922#discussion_r4146905044
########## task-sdk/docs/dynamic-task-mapping-vs-iteration.rst: ########## @@ -0,0 +1,446 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _sdk-dynamic-task-mapping-vs-iteration: + +Dynamic Task Mapping vs Iterable Tasks +====================================== + +.. versionadded:: 3.4.0 + +Airflow provides two complementary ways to process collections of data: + +- **Dynamic Task Mapping (DTM)** distributes work **across multiple workers**. + Each item becomes a separate Task Instance that can run on a different worker, + giving you horizontal scalability and per-item observability. + +- **Iterable Tasks (IT)** improves concurrency **within a single task**. + All items are processed inside one Task Instance on one worker, eliminating + scheduling overhead and — when combined with async operators — enabling true + I/O multiplexing through a shared event loop. + +In short: **DTM spreads load across workers; IT speeds up work within one worker.** + +While both approaches allow you to apply an operation over a collection, +they differ significantly in execution model, scheduler impact, and observability. +This page explains the trade-offs and when to use each. + +Real-World Motivation +--------------------- + +Consider a workflow that downloads ~17,000 XML files from an SFTP server and loads +them into a data warehouse. Community benchmarks demonstrate the dramatic performance +difference between the two approaches: + +.. list-table:: + :header-rows: 1 + + * - Approach + - Execution Time + * - Dynamic Task Mapping with mapped ``SFTPOperator`` + - 3 h 25 m + * - Sync ``@task`` with ``SFTPHook`` (sequential loop) + - 1 h 21 m + * - Async ``@task`` with ``SFTPHookAsync`` (concurrent loop) + - 8 m 29 s + * - Async ``@task`` with ``SFTPHookAsync`` and connection pooling + - 3 m 32 s + +The ~60× improvement stems from eliminating per-item scheduling overhead and +sharing a single event loop for concurrent I/O. This is the kind of workload +where IT excels: many small, I/O-bound operations processed within one task. + +Dynamic Task Mapping (DTM) +-------------------------- + +Dynamic Task Mapping allows you to expand a single task definition into multiple +Task Instances (TIs). + +For more details, see :ref:`dynamic task mapping <sdk-dynamic-task-mapping>`. + +Key characteristics: + +- Each item in the iterable creates a separate Task Instance. +- The scheduler is responsible for creating and managing all mapped tasks. +- Tasks can run in parallel across multiple worker slots. +- Fine-grained retry, logging, and observability per item. +- Well suited for workloads where each item should be independently scheduled and tracked. + +The following example fetches Pokémon data from a REST API. Each Pokémon becomes +a separate Task Instance, individually scheduled, retried, and visible in the UI: + +.. code-block:: python + + from datetime import datetime + + from airflow.providers.http.operators.http import HttpOperator + from airflow.sdk import DAG, task + + with DAG(dag_id="dtm-http-pokemon-example", start_date=datetime(2026, 1, 1)): + list_pokemon_task = HttpOperator( + task_id="list_pokemon", + http_conn_id="pokeapi", + method="GET", + endpoint="api/v2/pokemon?limit=100", + response_filter=lambda response: [ + pokemon["url"].replace("https://pokeapi.co/", "") for pokemon in response.json()["results"] + ], + log_response=False, + ) + + get_pokemon_task = HttpOperator.partial( + task_id="get_pokemon", + http_conn_id="pokeapi", + method="GET", + ).expand(endpoint=list_pokemon_task.output) + + list_pokemon_task >> get_pokemon_task + + +With 100 Pokémon the scheduler creates 100 Task Instances, each occupying +a worker slot. This is fine for small lists, but for thousands of items the +scheduler and database overhead becomes significant. + +Iterable Tasks (IT) +---------------------------- + +Iterable Tasks allows you to iterate over an iterable (typically an XCom result) +*within a single Task Instance*, applying an operator multiple times without creating +separate Task Instances. + +This means that iteration happens inside the task execution itself rather than at the +scheduler level. + +Key characteristics: + +- A single Task Instance processes all items in the iterable. +- No task expansion; the scheduler manages only one task. +- Lower scheduler overhead compared to DTM. +- Iterations share the same execution context (e.g., memory, event loop). +- Particularly well suited for async operators and high-throughput workloads. + +The same Pokémon fetching problem can be solved with IT. Here, a single Task +Instance processes all Pokémon concurrently using the sync +:class:`~airflow.providers.http.operators.http.HttpOperator`: + +.. code-block:: python + + from datetime import datetime + + from airflow.providers.http.operators.http import HttpOperator + from airflow.sdk import DAG, task + + with DAG(dag_id="it-http-pokemon-example", start_date=datetime(2026, 1, 1)): + list_pokemon_task = HttpOperator( + task_id="list_pokemon", + http_conn_id="pokeapi", + method="GET", + endpoint="api/v2/pokemon?limit=100", + response_filter=lambda response: [ + pokemon["url"].replace("https://pokeapi.co/", "") for pokemon in response.json()["results"] + ], + log_response=False, + ) + + get_pokemon_task = HttpOperator.partial( + task_id="get_pokemon", + http_conn_id="pokeapi", + method="GET", + ).iterate(endpoint=list_pokemon_task.output) + + list_pokemon_task >> get_pokemon_task + + +The scheduler only manages a single task. With sync tasks, iterations are +executed in a multi-threaded fashion, which eliminates scheduling overhead +and can speed up compute-bound workloads. However, for I/O-bound operations +like HTTP requests, multi-threading alone does not provide the same +performance benefits as async multiplexing — threads still block on each +request individually rather than sharing a single event loop. + +To truly **multiplex** I/O-bound operations, use an async task with +:class:`~airflow.providers.http.hooks.http.HttpAsyncHook`: + +.. code-block:: python + + from datetime import datetime + + from airflow.providers.http.hooks.http import HttpAsyncHook, HttpHook + from airflow.sdk import dag, task + + + @dag( + dag_id="it-async-http-pokemon-example", + start_date=datetime(2026, 1, 1), + ) + def it_async_http_pokemon_example(): + @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(), + ) + + + it_async_http_pokemon_example() + + +When ``iterate()`` is used with an async task, all iterations share the same +event loop, enabling true multiplexing of I/O-bound operations without any +manual concurrency management by the DAG author. For 5 Pokémon the +difference is negligible, but for hundreds or thousands of items the +concurrent approach is dramatically faster — see the +:ref:`benchmarks above <sdk-dynamic-task-mapping-vs-iteration>`. + +.. note:: + + ``multiple_outputs`` is ignored by ``iterate()``. Each iteration's return value is pushed + whole as ``return_value_<index>`` and the task's own return value is the lazy sequence over + 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. Review Comment: Thanks, f882327faa adds that to the note: a plain `threading.Thread` or `run_in_executor` gets the task's own context and un-indexed keys; use the `ti` passed in, or `asyncio.to_thread` / `copy_context().run`. There's a test pinning each of the three. Drafted-by: Claude Opus 5.5; reviewed by @dabla before posting -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
