Using Apache Airflow TaskFlow API to Simplify DAG Definition
Learn how the @task decorator turns Python functions into Airflow tasks with automatic XCom handling, see a concrete ETL example, and understand size limits and common pitfalls.
06 Feb 2026, 17:12 UTC

Why the TaskFlow API matters
Apache Airflow 2.0 introduced the @task decorator, which turns any plain Python function into an Airflow task. The decorator automatically pushes the function’s return value to XCom and makes that value available as an argument to downstream decorated functions. This removes the need for explicit set_upstream or >> bits and reduces boilerplate when you want to pass data between steps.
Worked example: a three‑step ETL pipeline
The following DAG defines extract, transform, and load steps as decorated functions. Dependencies are expressed simply by calling the functions in the order they should run.
from airflow import DAG
from airflow.decorators import task
from datetime import datetime
with DAG(
dag_id="etl_taskflow_example",
start_date=datetime(2026, 1, 1),
schedule_interval="@daily",
catchup=False,
) as dag:
@task
def extract() -> dict:
# Simulate pulling data from an API
return {"value": 42}
@task
def transform(raw: dict) -> dict:
# Example transformation: double the value
raw["value"] *= 2
return raw
@task
def load(processed: dict) -> None:
# Simulate writing to a destination
print(f"Loaded data: {processed}")
# Dependency chain: extract -> transform -> load
extracted = extract()
transformed = transform(extracted)
load(transformed)
What happens under the hood:
- Airflow serializes each decorated function and sends it to the worker.
- When the function executes, its return value is automatically pushed to XCom with a key that matches the task ID.
- The next decorated function receives that value as its argument because Airflow maps XCom pull to the function signature.
- The DAG edges are built from the order of the calls (
extract()→transform()→load()), so no explicit dependency statements are needed.
Limits and common pitfalls
XCom payload size
The return value of a @task function is stored in the metadata database. For MySQL the default column limits the payload to ~64 KB; PostgreSQL uses a TEXT column with no hard limit but very large values still affect performance. Returning large data structures (e.g., big Pandas DataFrames) should be avoided. Instead, write the data to an external system (S3, HDFS, a temporary file) and pass only a reference such as a file path or object key.
Function picklability
Because the function object is serialized with pickle and sent to the worker, it must be picklable. Avoid referencing non‑serializable objects inside the function, such as open file handles, database connections, or custom classes that lack a __reduce__ method. If you need a connection, create it inside the task function (or use Airflow hooks) rather than capturing an external object.
Mutable defaults and global state
Each task execution receives a fresh instance of the function. Mutable default arguments (def f(x=[]):) or modifications to module‑level globals are not shared between runs and can lead to confusing bugs. Prefer immutable defaults (None) and instantiate any needed containers inside the function.
Version requirement
The @task decorator is only available in Airflow 2.0+. Running the same DAG on an older version raises an ImportError. Verify the version before deployment:
airflow version # should output 2.0.0 or higher
Verification steps
- Trigger the DAG via the UI or
airflow dags trigger etl_taskflow_example. - Check the task logs:
airflow tasks log etl_taskflow_example extract 2026-10-05T00:00:00+00:00(adjust the execution date). You should see theprintstatement from theloadtask. - Optionally confirm XCom values:
airflow xcom pull -d etl_taskflow_example -t transformreturns the dictionary produced byextractafter transformation. - If you suspect a payload is too large, query the metadata DB directly (e.g.,
SELECT length(value) FROM xcom WHERE dag_id='etl_taskflow_example' AND task_id='extract';) and ensure it stays well below the limit.
When to stick with the classic operator style
If you need to pass very large objects, prefer writing them to external storage and using operators that handle staging (e.g., S3ToRedshiftOperator). For simple procedural steps that mainly orchestrate existing operators, the classic style may still be clearer. The TaskFlow API shines when the logic is naturally expressed as pure Python functions with modest data exchange.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.