Lesson 6 / 25

The TaskFlow API

Write DAGs with @dag and @task so Python functions become tasks and return values become XComs.

Plain functions as tasks

The TaskFlow API lets you write tasks as normal Python functions with the @task decorator and wire them by calling one with the output of another. Airflow stores the return value as an XCom and infers the dependency, so load(transform(extract())) creates the chain extract → transform → load. It is concise and the recommended style for Python-based tasks.

A small ETL with TaskFlow

The data is hard-coded to keep the example self-contained. In a real pipeline extract would call an API or database, and only small values would pass between tasks. I ran it with airflow dags test orders_etl 2026-10-01 on Airflow 3.0.6; the load task printed the total, 120 + 80 + 200 = 400.

from datetime import datetime

from airflow.sdk import dag, task


@dag(
    dag_id="orders_etl",
    start_date=datetime(2026, 9, 1),
    schedule="@daily",
    catchup=False,
    default_args={"retries": 2},
    tags=["demo"],
)
def orders_etl():
    @task
    def extract() -> list[dict]:
        return [{"id": 1, "amount": 120}, {"id": 2, "amount": 80}, {"id": 3, "amount": 200}]

    @task
    def transform(rows: list[dict]) -> int:
        return sum(r["amount"] for r in rows)

    @task
    def load(total: int) -> None:
        print(f"Total revenue loaded: {total}")

    load(transform(extract()))


orders_etl()

Output:

Total revenue loaded: 400

Type-hint your task functions

Return type hints document what flows between tasks and help catch mistakes early. Keep the values small and serialisable (numbers, strings, small lists and dicts).

Quick check: In TaskFlow, how does data pass from one @task to the next?

  • It cannot pass data
  • Through a global variable
  • Through email
  • As the return value, stored as an XCom
Answer

As the return value, stored as an XCom — Return values become XComs and Airflow wires the dependency automatically.