With the TaskFlow API, @task turns a function into a task. Calling it inside a DAG adds the task and returns an XComArg, a promise of its return value; passing that promise on hands over the data and creates the dependency:
@task
def extract(ds=None) -> str:
path, n = p.extract_orders(ds, f"{DATA}/source/orders.jsonl", f"{DATA}/landing")
print(f"extracted {n} orders for {ds}")
return path # becomes an XCom for load()
@task
def load(path: str, ds=None) -> dict:
orders, lines = p.load_orders(ds, path, booknest_conn())
return {"orders": orders, "lines": lines}
load(extract())Parameters named after context variables (ds, logical_date, ti) are filled in at run time. Return values travel as XComs: the supervisor posts them to the API server, which stores them in the xcom table:
task_id | key | value
-----------------+--------------+--------------------------------------------------------------
-----------
ingest.extract | return_value | "/opt/airflow/data/landing/orders/date=2026-06-29/orders.json
l"
ingest.load | lines | 261
ingest.load | orders | 192
ingest.load | return_value | {"lines": 261, "orders": 192}
model.check | return_value | "261 lines, gross 6622.95"
model.dbt_build | return_value | "21:37:55 Done. PASS=6 WARN=0 ERROR=0 SKIP=0 NO-OP=0
REUSED=0 TOTAL=6"
publish | return_value | 6Because load is annotated -> dict, TaskFlow set multiple_outputs=True and stored each key too; the BashOperator pushed its last output line. XComs are metadata for paths, IDs and counts: for real data, pass a storage location or use the common-io provider's object-storage XCom backend.