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.