Menu

Airflow course · Lesson 1 of 5

Airflow DAG Fundamentals, TaskFlow and Dynamic DAGs

Build Airflow 3 DAGs from first principles: tasks and dependencies, the TaskFlow API, params and Jinja templating, dynamic DAGs and dynamic task mapping.

  • Beginner
  • 26 min read
  • Updated Oct 2026
On this page
  1. Running the examples
  2. DAG fundamentals
  3. What it is
  4. How it works
  5. Pitfalls
  6. In interviews
  7. Tasks and dependencies
  8. What it is
  9. How it works
  10. Pitfalls
  11. In interviews
  12. The TaskFlow API
  13. What it is
  14. How it works
  15. Pitfalls
  16. In interviews
  17. DAG params and templating with Jinja
  18. What it is
  19. How it works
  20. Pitfalls
  21. In interviews
  22. Dynamic DAG generation
  23. What it is
  24. How it works
  25. Pitfalls
  26. In interviews
  27. Dynamic task mapping
  28. What it is
  29. How it works
  30. Pitfalls
  31. In interviews
  32. DAG authoring best practices
  33. In interviews
  34. Practice questions
  35. Key takeaways

A DAG is how you tell Airflow what to run, in what order and on what schedule. Every other Airflow skill (sensors, retries, backfills, scaling) builds on writing DAGs that are correct, cheap to parse and safe to rerun. This lesson covers the building blocks in Airflow 3: tasks and dependencies, the TaskFlow API, params and templates, and the two ways of making pipelines dynamic.

Running the examples

The examples use Airflow 3.3. To follow along, install Airflow in a virtual environment (pip install apache-airflow) and create the metadata database once with airflow db migrate.

In Airflow 3, dag.test() runs one DAG run in a single process without a scheduler, but it first looks the DAG up in the metadata database. Normally that works because you call it from a file in your DAG folder (python dags/my_dag.py with if __name__ == "__main__": dag.test() at the bottom), and Airflow parses the folder first. The snippets on this page are not files in a DAG folder, so this small helper registers the in-memory DAG, runs it, and returns the run state, each task’s state and each task’s return value.

import contextlib, io, tempfile
from airflow.dag_processing.bundles.manager import DagBundlesManager
from airflow.dag_processing.dagbag import DagBag, sync_bag_to_db
from airflow.models.xcom import XComModel
from airflow.utils.session import create_session

def run_dag(dag, **test_kwargs):
    """Register an in-memory DAG, run it once with dag.test(), return (run state, task states, return values)."""
    with contextlib.redirect_stdout(io.StringIO()):          # hide Airflow's own log lines
        DagBundlesManager().sync_bundles_to_db()
        bag = DagBag(dag_folder=tempfile.mkdtemp())
        bag.dags[dag.dag_id] = dag
        sync_bag_to_db(bag, "dags-folder", None)
        dr = dag.test(**test_kwargs)
    states, values = {}, {}
    with create_session() as session:
        for ti in dr.get_task_instances(session=session):
            key = ti.task_id if ti.map_index < 0 else f"{ti.task_id}[{ti.map_index}]"
            states[key] = str(ti.state)
        for x in session.query(XComModel).filter_by(dag_id=dag.dag_id, run_id=dr.run_id, key="return_value"):
            key = x.task_id if x.map_index < 0 else f"{x.task_id}[{x.map_index}]"
            values[key] = XComModel.deserialize_value(x)
    return str(dr.state), dict(sorted(states.items())), dict(sorted(values.items()))

You do not need this helper in a real project. Use python dags/my_dag.py or airflow dags test <dag_id> <date> instead.

DAG fundamentals

What it is

A DAG (directed acyclic graph) is a collection of tasks plus the dependencies between them. Directed means each dependency has a direction (extract before load). Acyclic means there are no loops: a task can never end up waiting on itself. The DAG also carries scheduling settings: when it starts, how often it runs, and whether missed runs are created.

Four words are used constantly, and interviewers expect you to keep them apart:

Term Meaning
DAG The definition: tasks, dependencies, schedule. Lives in a Python file.
Task One node in the DAG, created from an operator or a @task function.
DAG run One execution of the whole DAG, identified by a run_id and usually a logical_date.
Task instance One execution of one task inside one DAG run (plus a map_index for mapped tasks and a try_number for retries).

How it works

You write a Python file that builds a DAG object and put it in a DAG bundle (by default the dags_folder). The DAG processor imports that file repeatedly, serialises each DAG it finds to JSON and stores it in the metadata database. The scheduler and the UI work from the serialised copy, and workers re-import your file only when they run a task. Two consequences follow: the file is executed many times, so its top-level code must be cheap; and a DAG that fails to import shows up as an import error, not as a failed run.

There are three equivalent ways to declare a DAG:

from datetime import datetime
from airflow.sdk import DAG, dag
from airflow.providers.standard.operators.empty import EmptyOperator   # Airflow 2: airflow.operators.empty

# 1. Context manager: operators created inside the block join the DAG
with DAG(dag_id="ctx_style", schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False) as ctx_dag:
    EmptyOperator(task_id="start") >> EmptyOperator(task_id="end")

# 2. Decorator: the function body builds the DAG, calling it returns the DAG object
@dag(schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False, tags=["demo"])
def decorator_style():
    EmptyOperator(task_id="start") >> EmptyOperator(task_id="end")

deco_dag = decorator_style()   # must be called, or no DAG is registered

# 3. Plain constructor, passing dag= to each operator
plain_dag = DAG(dag_id="plain_style", schedule=None, start_date=datetime(2026, 1, 1))
EmptyOperator(task_id="start", dag=plain_dag) >> EmptyOperator(task_id="end", dag=plain_dag)

for d in (ctx_dag, deco_dag, plain_dag):
    print(d.dag_id, d.task_ids, "catchup:", d.catchup, "max_active_runs:", d.max_active_runs)
ctx_style ['start', 'end'] catchup: False max_active_runs: 16
decorator_style ['start', 'end'] catchup: False max_active_runs: 16
plain_style ['start', 'end'] catchup: False max_active_runs: 16

The @dag decorator uses the function name as the dag_id unless you pass one. In Airflow 3 the public authoring interface is airflow.sdk (in Airflow 2 you imported DAG from airflow or airflow.models and dag/task from airflow.decorators), and the everyday operators live in the apache-airflow-providers-standard package.

The settings you set on almost every DAG:

Argument What it controls Note
dag_id Unique name Changing it creates a new DAG with no history.
schedule When runs are created Cron string, preset such as "@daily", timedelta, timetable, asset expression, or None for manual only. Airflow 2.4 renamed schedule_interval to schedule; Airflow 3 removed the old name.
start_date Earliest date the schedule considers Use a fixed, timezone-aware date, never datetime.now().
catchup Whether missed intervals since start_date are run Defaults to False in Airflow 3 (True in Airflow 2). Set it explicitly.
default_args Defaults passed to every task (retries, owner, timeouts) A task’s own argument wins.
max_active_runs, max_active_tasks Concurrency limits for this DAG Defaults come from config (16 each in a default 3.3 install).
tags, doc_md, owner_links UI metadata Cheap and useful for discovery.

Pitfalls

  • Forgetting to call a @dag-decorated function at module level. The file imports cleanly but no DAG appears.
  • Reusing a dag_id in two files. Only one wins, and which one depends on parse order.
  • A dynamic start_date such as datetime.now(). The schedule keeps moving, so runs may never be created.
  • Treating the DAG file as a script that “runs the pipeline”. It only describes the pipeline; code at the top level runs every time the file is parsed.

In interviews

“What is a DAG and why must it be acyclic?” is a warm-up. A strong answer separates DAG, DAG run and task instance, explains that the DAG processor parses files continuously and serialises them for the scheduler, and mentions why top-level code must be light. Knowing that Airflow 3 moved the authoring API to airflow.sdk and the standard operators into a provider shows you have used a current version.

Tasks and dependencies

What it is

A task is a unit of work: run a SQL statement, call an API, submit a Spark job. Each task is an instance of an operator (a template for a kind of work) or a function decorated with @task. Dependencies say which tasks must finish before others start. By default a task runs only when all its direct upstream tasks succeeded (the all_success trigger rule).

How it works

  • a >> b means “a before b” (b << a is the same). a >> [b, c] fans out; [b, c] >> d fans in.
  • chain(a, b, c) links a sequence. With two lists of equal length it links them pairwise.
  • cross_downstream([a, b], [c, d]) makes every task in the first list upstream of every task in the second. You cannot write [a, b] >> [c, d] because Python cannot apply >> between two lists.
  • A TaskGroup groups tasks visually and prefixes their ids (group_id.task_id). Dependencies on the group apply to its first and last tasks. Task groups replaced SubDAGs, which were removed in Airflow 3.
from airflow.sdk import DAG, TaskGroup, chain, cross_downstream
from airflow.providers.standard.operators.empty import EmptyOperator

with DAG("deps_demo", schedule=None, start_date=datetime(2026, 1, 1)) as deps_dag:
    start = EmptyOperator(task_id="start")
    extract_a, extract_b = EmptyOperator(task_id="extract_a"), EmptyOperator(task_id="extract_b")
    check, stage = EmptyOperator(task_id="check"), EmptyOperator(task_id="stage")
    with TaskGroup("publish") as publish:
        load = EmptyOperator(task_id="load")
        notify = EmptyOperator(task_id="notify")
        load >> notify

    start >> [extract_a, extract_b]                          # fan out
    cross_downstream([extract_a, extract_b], [check, stage]) # every extract before check and stage
    chain([check, stage], publish)                           # both before the group

print(deps_dag.task_ids)
for t in deps_dag.tasks:
    print(f"{t.task_id:16} -> {sorted(t.downstream_task_ids)}")
['start', 'extract_a', 'extract_b', 'check', 'stage', 'publish.load', 'publish.notify']
start            -> ['extract_a', 'extract_b']
extract_a        -> ['check', 'stage']
extract_b        -> ['check', 'stage']
check            -> ['publish.load']
stage            -> ['publish.load']
publish.load     -> ['publish.notify']
publish.notify   -> []

Building a cycle does not fail on the spot. The check runs when the DAG file is parsed, and the DAG processor then records an import error instead of loading the DAG. You can run the same check yourself, which is what a DAG integrity test does:

with DAG("cycle_demo", schedule=None, start_date=datetime(2026, 1, 1)) as cycle_dag:
    a, b = EmptyOperator(task_id="a"), EmptyOperator(task_id="b")
    a >> b >> a          # no error yet

try:
    cycle_dag.check_cycle()
except Exception as exc:
    print(type(exc).__name__, "-", exc)
AirflowDagCycleException - Cycle detected in Dag: cycle_demo. Faulty task: b

Pitfalls

  • Setting dependencies inside a loop and accidentally linking every task to every other task. Print downstream_task_ids or look at the graph view.
  • Assuming tasks run in the order they appear in the file. Only declared dependencies matter; independent tasks can run in any order and in parallel.
  • Many tiny tasks. Each task instance has scheduling overhead (often seconds), so a DAG of hundreds of trivial tasks is slow. Group trivial steps into one task, and split only where a retry boundary or parallelism helps.
  • Using a dependency to pass data. Dependencies order work; data moves through XCom (small values) or storage (anything large).

In interviews

Expect “How do you make task C wait for both A and B?” ([a, b] >> c) and “How do you link two lists?” (cross_downstream, or chain for pairwise). For a stronger answer, mention task groups as the replacement for SubDAGs, and that trigger rules change the default “all upstream succeeded” condition.

The TaskFlow API

What it is

The TaskFlow API lets you write tasks as plain Python functions decorated with @task. Calling a decorated function inside a DAG does not run it. It creates a task and returns an XComArg, a placeholder for the value the task will return at run time. Passing that placeholder to another task both moves the data (through XCom) and creates the dependency.

How it works

from airflow.sdk import dag, task

@dag(schedule=None, start_date=datetime(2026, 1, 1), catchup=False)
def orders_taskflow():
    @task
    def extract() -> list[dict]:
        return [{"id": 1, "amount": 40.0}, {"id": 2, "amount": 15.5}, {"id": 3, "amount": -1.0}]

    @task(multiple_outputs=True)          # each dict key becomes its own XCom
    def transform(rows: list[dict]) -> dict:
        good = [r for r in rows if r["amount"] >= 0]
        return {"rows": good, "rejected": len(rows) - len(good)}

    @task
    def load(rows: list[dict], rejected: int) -> str:
        return f"loaded {len(rows)} rows, rejected {rejected}"

    cleaned = transform(extract())
    load(cleaned["rows"], cleaned["rejected"])

taskflow_dag = orders_taskflow()
print(taskflow_dag.task_ids)
print(run_dag(taskflow_dag))
['extract', 'transform', 'load']
('success', {'extract': 'success', 'load': 'success', 'transform': 'success'}, {'extract': [{'id': 1, 'amount': 40.0}, {'id': 2, 'amount': 15.5}, {'id': 3, 'amount': -1.0}], 'load': 'loaded 2 rows, rejected 1', 'transform': {'rows': [{'id': 1, 'amount': 40.0}, {'id': 2, 'amount': 15.5}], 'rejected': 1}})

Things to notice:

  • Dependencies extract >> transform >> load came from passing return values, with no >> written.
  • With multiple_outputs=True (or a dict return type hint), each key is stored as a separate XCom, so a downstream task can pull just one key.
  • The values travelled through XCom, which lives in the metadata database by default. That is fine for a few kilobytes and wrong for a data frame. Pass a path or table name instead.

You can mix TaskFlow with classic operators. The .output property of an operator is its XComArg, and >> still works between both kinds:

from airflow.providers.standard.operators.bash import BashOperator

@dag(schedule=None, start_date=datetime(2026, 1, 1))
def mixed_styles():
    today = BashOperator(task_id="today", bash_command="echo 2026-03-01")   # pushes its last stdout line

    @task
    def build_path(date_str: str) -> str:
        return f"s3://lake/orders/dt={date_str}/"

    path = build_path(today.output)

    @task
    def audit():
        return "audit done"

    path >> audit()

print(run_dag(mixed_styles())[2])
{'audit': 'audit done', 'build_path': 's3://lake/orders/dt=2026-03-01/', 'today': '2026-03-01'}

Because a decorated task wraps a normal function, you can unit test the logic without Airflow at all through .function:

from airflow.sdk import task

@task
def dedupe(ids: list[int]) -> list[int]:
    return sorted(set(ids))

assert dedupe.function([3, 1, 3, 2]) == [1, 2, 3]
print("dedupe logic ok")
dedupe logic ok

Other decorators follow the same pattern: @task.bash (the function returns the command to run), @task.virtualenv and @task.external_python (run in a different Python environment), @task.branch (choose a path, covered in the next lesson), @task.sensor, and provider decorators such as @task.docker and @task.kubernetes.

Pitfalls

  • Returning large objects. Every return value is serialised into XCom. Return a reference.
  • Calling a task function at the top level outside a DAG, expecting it to compute something. It only builds a task.
  • Doing work in the DAG function body (outside any @task). That code runs at parse time, every time the file is parsed.
  • Forgetting that each task may run on a different worker in a different process: globals set in one task are not visible in another.

In interviews

“TaskFlow vs classic operators?” A good answer: TaskFlow is cleaner for Python logic, infers dependencies from data flow and handles XCom automatically; classic operators are still the way to use provider integrations (SQL, Spark, Kubernetes) and the two mix freely. Mention multiple_outputs, the XCom size caveat, and that the Airflow 2 import was from airflow.decorators import dag, task.

DAG params and templating with Jinja

What it is

Many operator arguments are templated: before a task runs, Airflow renders them with Jinja using the task’s context (the run’s dates, params, the task instance and more). Templates are how a task learns which slice of data it is responsible for, which is what makes reruns and backfills safe. Params are run-level inputs with defaults and validation that users can override when they trigger a run.

How it works

Each operator lists its templated arguments in template_fields (for example bash_command and env on BashOperator, sql on SQL operators, op_args, op_kwargs and templates_dict on PythonOperator). The most used template variables:

Variable Value for a daily run covering 1 March 2026 Note
{{ ds }} 2026-03-01 Logical date as YYYY-MM-DD.
{{ ds_nodash }} 20260301 Handy for file names.
{{ logical_date }} 2026-03-01 00:00:00+00:00 Airflow 2’s execution_date was removed in Airflow 3.
{{ data_interval_start }} / {{ data_interval_end }} start and end of the interval For a cron schedule in Airflow 3 they are equal unless you opt into data intervals (see the scheduling lesson).
{{ params.name }} a param value Run conf overrides the default.
{{ macros.ds_add(ds, -1) }} 2026-02-28 Built-in macros for date arithmetic.
{{ var.value.key }}, {{ conn.my_conn.host }} a Variable or Connection field Read at render time, not parse time.

Params are declared with Param, which validates values against JSON Schema:

import pendulum
from airflow.sdk import DAG, Param
from airflow.providers.standard.operators.bash import BashOperator

with DAG(
    "templating_demo",
    schedule="@daily",
    start_date=datetime(2026, 1, 1),
    catchup=False,
    params={
        "region": Param("eu", type="string", enum=["eu", "us"]),
        "limit": Param(100, type="integer", minimum=1),
    },
) as tmpl_dag:
    BashOperator(
        task_id="show",
        bash_command=(
            "echo ds={{ ds }} yesterday={{ macros.ds_add(ds, -1) }} "
            "region={{ params.region }} limit={{ params.limit }}"
        ),
    )

march_first = pendulum.datetime(2026, 3, 1, tz="UTC")
print(run_dag(tmpl_dag, logical_date=march_first)[2])
print(run_dag(tmpl_dag, logical_date=march_first, run_conf={"region": "us", "limit": 5})[2])
{'show': 'ds=2026-03-01 yesterday=2026-02-28 region=eu limit=100'}
{'show': 'ds=2026-03-01 yesterday=2026-02-28 region=us limit=5'}

An invalid value is rejected before any task runs:

try:
    run_dag(tmpl_dag, logical_date=march_first, run_conf={"region": "apac"})
except Exception as exc:
    print(type(exc).__name__, "-", str(exc).splitlines()[0])
ParamValidationError - Invalid input for param region: 'apac' is not one of ['eu', 'us']

And the Airflow 2 variable is simply gone. The task fails with jinja2.exceptions.UndefinedError: 'execution_date' is undefined in its log:

with DAG("old_template", schedule=None, start_date=datetime(2026, 1, 1)) as old_dag:
    BashOperator(task_id="old", bash_command="echo {{ execution_date }}")

print(run_dag(old_dag, logical_date=march_first)[1])
{'old': 'failed'}

Inside a TaskFlow function you do not need Jinja: add a parameter named after a context key (def load(ds=None, params=None)) or call get_current_context().

Pitfalls

  • Templating a field that is not in template_fields. The string is passed through literally, braces and all. Check the operator’s documentation or Operator.template_fields.
  • Templates are rendered only when the task runs, never at parse time. print("{{ ds }}") in top-level code prints the braces.
  • datetime.now() instead of {{ ds }} or data_interval_*. A rerun of last Tuesday would then process today’s data.
  • Jinja renders to strings. To get a native list or dict in a templated argument, set render_template_as_native_obj=True on the DAG.
  • A templated string that ends in .sql or .sh is treated as a path to a template file, which surprises people passing a literal command.

In interviews

Expect “How do you make a task process the right day’s data?” (template ds or data_interval_start/end into the query, never use the wall clock) and “What replaced execution_date?” (logical_date, plus the data interval variables). Mention that params are validated with JSON Schema and can be overridden by trigger conf.

Dynamic DAG generation

What it is

Dynamic DAG generation means one Python file creates several DAGs, usually by looping over configuration: one DAG per source system, per customer or per table. The structure is fixed at parse time.

How it works

The DAG processor discovers DAGs by finding DAG objects created when the file is imported. With the with DAG(...) or @dag style, creating the DAG inside the loop is enough (auto_register=True is the default):

SOURCES = {
    "crm": {"schedule": "@hourly", "tables": ["accounts", "contacts"]},
    "billing": {"schedule": "@daily", "tables": ["invoices"]},
}

generated = {}
for source, cfg in SOURCES.items():
    with DAG(
        dag_id=f"ingest_{source}",
        schedule=cfg["schedule"],
        start_date=datetime(2026, 1, 1),
        catchup=False,
        tags=["generated", source],
    ) as generated_dag:
        done = EmptyOperator(task_id="done")
        for table in cfg["tables"]:
            BashOperator(task_id=f"copy_{table}", bash_command=f"echo copy {source}.{table}") >> done
    generated[generated_dag.dag_id] = generated_dag

for dag_id, d in generated.items():
    print(dag_id, type(d.timetable).__name__, d.task_ids)
ingest_crm CronTriggerTimetable ['done', 'copy_accounts', 'copy_contacts']
ingest_billing CronTriggerTimetable ['done', 'copy_invoices']

Each DAG gets its own schedule, and a cron preset such as @hourly becomes a CronTriggerTimetable in Airflow 3 (the scheduling lesson explains why that matters). Common sources of configuration are a Python dict or a YAML/JSON file next to the DAG file. Reading a local file at parse time is fine. Calling an API or querying a database at parse time is not: it runs every time the file is parsed (every 30 seconds by default per file, controlled by [dag_processor] min_file_process_interval), slows parsing for every DAG and can time out ([core] dagbag_import_timeout, 30 seconds by default).

Pitfalls

  • Generating from a remote source at parse time. Generate the config file in CI or by a separate job, and read it locally.
  • Removing an entry from the config makes its DAG disappear from the UI, while its history stays in the database.
  • Hundreds of DAGs in one file slow down parsing of that file. Split by domain.
  • Generated dag_ids must be stable, since a renamed DAG loses its run history.

In interviews

“You have 200 tables to ingest with the same logic. One DAG or many?” Talk about the trade-off: one DAG per source gives independent schedules, retries and ownership; one DAG with many tasks or mapped tasks is easier to monitor as a unit. Say that generation config must be static and local at parse time, and mention dynamic task mapping when the list is only known at run time.

Dynamic task mapping

What it is

Dynamic task mapping creates a variable number of copies of a task at run time, one per input element, like a map in functional programming. Unlike dynamic DAG generation, the list does not need to exist when the file is parsed: it can be the output of an upstream task (for example, “the files that arrived today”).

How it works

  • .expand(arg=iterable) creates one mapped task instance per element. Each gets a map_index (0, 1, 2…).
  • .partial(...) fixes the arguments that are the same for every copy.
  • Expanding over two arguments gives the cross product.
  • .expand_kwargs(list_of_dicts) maps over whole argument sets, so arguments vary together rather than as a cross product.
  • A downstream task that receives the mapped output gets a lazy sequence of all results (a “reduce” step).
  • Classic operators map too: BashOperator.partial(task_id=...).expand(bash_command=[...]).
  • map_index_template gives each copy a readable label in the UI instead of a bare index.
@dag(schedule=None, start_date=datetime(2026, 1, 1))
def mapping_demo():
    @task
    def list_files() -> list[str]:                  # known only at run time
        return ["orders_eu.csv", "orders_us.csv", "orders_apac.csv"]

    @task(map_index_template="{{ file_name }}")
    def count_rows(file_name: str, min_rows: int) -> dict:
        from airflow.sdk import get_current_context
        get_current_context()["file_name"] = file_name   # used by map_index_template
        return {"file": file_name, "rows": len(file_name) * 10}

    @task
    def total(results) -> int:                      # receives every mapped result
        return sum(r["rows"] for r in results)

    total(count_rows.partial(min_rows=5).expand(file_name=list_files()))

    @task
    def add(x: int, y: int) -> int:
        return x + y

    add.expand(x=[1, 2], y=[10, 20])                # cross product: 4 task instances

    @task
    def copy(src: str, dest: str) -> str:
        return f"{src}->{dest}"

    copy.expand_kwargs([{"src": "a", "dest": "x"}, {"src": "b", "dest": "y"}])   # 2 instances

    BashOperator.partial(task_id="echo").expand(bash_command=["echo a", "echo b"])

state, states, values = run_dag(mapping_demo())
print(state)
for key, value in values.items():
    print(f"{key:14} {value}")
success
add[0]         11
add[1]         21
add[2]         12
add[3]         22
copy[0]        a->x
copy[1]        b->y
count_rows[0]  {'file': 'orders_eu.csv', 'rows': 130}
count_rows[1]  {'file': 'orders_us.csv', 'rows': 130}
count_rows[2]  {'file': 'orders_apac.csv', 'rows': 150}
echo[0]        a
echo[1]        b
list_files     ['orders_eu.csv', 'orders_us.csv', 'orders_apac.csv']
total          410

add ran four times (1+10, 1+20, 2+10, 2+20), copy twice, and total added up the three mapped results. In the UI the count_rows copies are labelled by file name thanks to map_index_template.

What happens with an empty list? The mapped task gets no instances and is marked skipped, and its downstream tasks are skipped too under the default trigger rule:

@dag(schedule=None, start_date=datetime(2026, 1, 1))
def empty_mapping():
    @task
    def nothing_new() -> list[str]:
        return []

    @task
    def process(f: str) -> str:
        return f

    @task
    def summarise(results) -> int:
        return len(list(results))

    summarise(process.expand(f=nothing_new()))

print(run_dag(empty_mapping())[1])
{'nothing_new': 'success', 'process': 'skipped', 'summarise': 'skipped'}

Limits worth knowing: [core] max_map_length (1024 by default) caps how many copies one expand can create, and the mapped instances count against the normal concurrency limits; use max_active_tis_per_dag on the task to throttle them.

Pitfalls

  • Mapping over something huge (one task per row). Map over batches or partitions, not records.
  • Expanding over two arguments when you wanted pairs. Use expand_kwargs, or zip the upstream outputs (a.zip(b)).
  • Only keyword arguments can be mapped, and task_id, pool and similar scheduling arguments go in partial, not expand.
  • Forgetting that an empty input skips everything downstream. If a summary task must always run, give it a trigger rule such as all_done or none_failed.

In interviews

“How would you process a variable number of files that arrive each day?” is the classic prompt. Answer with a listing task that returns paths, a mapped processing task (partial for constants, expand for the paths), and a reduce task; mention max_map_length, throttling with max_active_tis_per_dag, and the empty-list skip behaviour. Contrast it with dynamic DAG generation, which is fixed at parse time.

DAG authoring best practices

These are the habits that separate a DAG that works on your laptop from one that runs reliably for years. The best practices lesson at the end of this course covers idempotency, top-level code and testing in depth.

  1. Keep the DAG file declarative and cheap to import. No database queries, API calls or heavy imports at module level. Put heavy imports inside the task function.
  2. Make every task idempotent and driven by the run’s dates. Use ds or data_interval_start/end (or logical_date), never the wall clock, and overwrite or merge rather than append.
  3. Orchestrate, do not process. Push heavy work to Spark, the warehouse or a container; workers should mostly wait.
  4. Pass references, not data. XCom carries file paths, table names and small metadata.
  5. Set scheduling arguments explicitly: a fixed start_date, catchup, max_active_runs, and retries/retry_delay in default_args. Defaults changed between Airflow 2 and 3.
  6. Size tasks around retry boundaries. One task per step you might want to rerun alone; not one per SQL statement, not one giant task.
  7. Name and document. Stable dag_ids, meaningful task_ids, tags, doc_md and an owner.
  8. Use timeouts. execution_timeout on tasks and dagrun_timeout on the DAG stop runs that hang forever.

In interviews

When asked “What makes a good DAG?”, lead with idempotency and light top-level code, then mention explicit scheduling settings and keeping heavy compute out of the workers. Concrete failure stories (duplicates after a retry, a slow scheduler caused by an API call in a DAG file) make the answer convincing.

Practice questions

What is the difference between a DAG run and a task instance?

A DAG run is one execution of the whole DAG, identified by a run_id and (for scheduled runs) a logical date and data interval. A task instance is one task inside one DAG run; it has its own state, try number and, for mapped tasks, a map index. Retrying a task creates a new try of the same task instance, not a new DAG run.

Why is calling an API at the top level of a DAG file a problem?

The DAG processor imports every DAG file repeatedly (every 30 seconds per file by default), and workers import it again before each task. A top-level API call therefore runs constantly, slows parsing for every DAG, can hit dagbag_import_timeout and turn the DAG into an import error, and may hammer the API. Move the call into a task or generate static config ahead of time.

How do you make d run after both b and c, which both run after a?

a >> [b, c] >> d. Lists work on one side of >>. For list-to-list dependencies use cross_downstream([b, c], [d, e]), or chain for pairwise links.

A task must process “yesterday’s” files. How do you get the date?

Template the run’s own dates: {{ ds }}, {{ data_interval_start }}, {{ logical_date }} or macros.ds_add(ds, -1), depending on the schedule’s semantics. In a TaskFlow task, accept ds or data_interval_start as a parameter or call get_current_context(). Never use datetime.now(), because a rerun or backfill must process the same slice of data it processed the first time.

Dynamic DAG generation or dynamic task mapping: which would you use to process every file that landed in a bucket today?

Dynamic task mapping. The set of files is only known at run time, so a listing task returns the paths and a mapped task processes each one, followed by an optional reduce task. Dynamic DAG generation creates DAGs from configuration that exists at parse time, such as a list of source systems.

What does multiple_outputs=True do on a @task?

It stores each key of the returned dict as its own XCom (plus the whole dict), so downstream tasks can depend on a single key, for example result["rows"]. A return type hint of dict turns it on automatically.

What happened to execution_date, SubDAGs and schedule_interval in Airflow 3?

execution_date was removed from the context; use logical_date (and the data_interval_* values). SubDAGs were removed; use TaskGroups for grouping, or separate DAGs connected by assets or TriggerDagRunOperator. schedule_interval (and timetable) were replaced by the single schedule argument, which was introduced in Airflow 2.4.

Key takeaways

  • A DAG file only describes a pipeline. The DAG processor imports it repeatedly, so top-level code must be cheap and deterministic.
  • Dependencies come from >>, chain, cross_downstream, task groups, or from passing TaskFlow return values.
  • TaskFlow turns functions into tasks and moves return values through XCom, so pass references rather than data.
  • Templates and params tell each run which slice of data to process. logical_date replaced execution_date in Airflow 3.
  • Generate DAGs from static, local config at parse time; use dynamic task mapping (expand, partial, expand_kwargs) when the work list is only known at run time.
  • Set start_date, catchup, retries and timeouts explicitly, because several defaults changed between Airflow 2 and 3.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Examples run on Apache Airflow 3.3.2 (Task SDK 1.3.2, apache-airflow-providers-standard 1.20.0) with Python 3.11 and a SQLite metadata database created by `airflow db migrate`. Airflow 2 names are given where they differ.

Progress is saved in this browser only. No account needed.

Search
Filter by type