The TaskFlow API and XComs

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:

booknest_daily.py (excerpt): data flow is the dependencyPython
@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:

Output of 15
     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 | 6

Because 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.