Airflow courseLesson 7 of 8
Airflow course · Lesson 7 of 8
Airflow Best Practices: Idempotent Tasks, Testing and CI/CD
Write Airflow DAGs that are safe to rerun and cheap to parse, then test them with DagBag integrity tests, unit tests and dag.test(), and ship them through CI/CD.
On this page
Most Airflow incidents are not Airflow bugs. They are tasks that duplicate data when retried, DAG files that call an API every time they are parsed, and changes that reach production without anyone importing the file first. This lesson covers the habits that prevent them: idempotent tasks, light top-level code, a test suite for DAGs, and a CI/CD pipeline that runs it on every change.
Running the examples
The examples use Airflow 3.3 with a metadata database created by airflow db migrate, and the standard library’s sqlite3 as a stand-in warehouse. The helper writes small DAG files into a temporary DAG folder so the tests below can load real files, the way the DAG processor does.
import contextlib, io, os, sqlite3, tempfile, textwrap, time
from datetime import datetime
import pendulum
from airflow.dag_processing.bundles.manager import DagBundlesManager
from airflow.dag_processing.dagbag import DagBag, sync_bag_to_db
from airflow.utils.session import create_session
DAGS_FOLDER = tempfile.mkdtemp(prefix="dags_")
WAREHOUSE = os.path.join(tempfile.mkdtemp(), "warehouse.db")
os.environ["DEMO_WAREHOUSE"] = WAREHOUSE # read by the DAG files at run time, not parse time
def write_dag(file_name: str, source: str) -> None:
with open(os.path.join(DAGS_FOLDER, file_name), "w") as fh:
fh.write(textwrap.dedent(source))
def run_dag(dag, **test_kwargs):
"""Register a DAG (normally the DAG processor's job), run it once with dag.test(), return task states."""
with contextlib.redirect_stdout(io.StringIO()):
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)
with create_session() as session:
return str(dr.state), {ti.task_id: str(ti.state) for ti in dr.get_task_instances(session=session)}
Idempotent tasks
What it is
A task is idempotent when running it once or many times for the same logical date leaves the target in the same final state. Retries, clears, backfills and at-least-once delivery all rerun tasks, so idempotency is what makes them safe. Without it every rerun is a potential duplicate or corruption.
How it works
Three patterns cover almost every case:
| Pattern | How | Good for |
|---|---|---|
| Overwrite the partition | delete the run’s slice and insert it again in one transaction, or INSERT OVERWRITE / partition replace |
daily or hourly partitions in a warehouse or lake |
| Merge on a key | MERGE / INSERT ... ON CONFLICT DO UPDATE with a natural or business key |
dimension tables, CDC, late-arriving updates |
| Write then swap | write to a temporary table or path, then rename or swap atomically | full refreshes and file outputs |
All three need the task to be parameterised by the run’s dates (ds, data_interval_start, logical_date), never by the wall clock, so that a rerun of 1 March touches only 1 March.
Here are two versions of the same daily load. One appends, one replaces its partition in a single transaction. Each DAG runs twice for the same day, as a retry or a clear would:
write_dag("orders_load.py", '''
import os, sqlite3
from datetime import datetime
from airflow.sdk import dag, task
def source_rows(ds):
return [(ds, "o-1", 40.0), (ds, "o-2", 15.5), (ds, "o-3", 9.0)] # pretend extract for one day
@dag(schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False, tags=["orders"],
default_args={"owner": "sales-data", "retries": 2})
def orders_append():
@task
def load(ds=None):
with sqlite3.connect(os.environ["DEMO_WAREHOUSE"]) as con:
con.execute("CREATE TABLE IF NOT EXISTS orders_a (ds TEXT, order_id TEXT, amount REAL)")
con.executemany("INSERT INTO orders_a VALUES (?, ?, ?)", source_rows(ds)) # appends
load()
@dag(schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False, tags=["orders"],
default_args={"owner": "sales-data", "retries": 2})
def orders_overwrite():
@task
def load(ds=None):
with sqlite3.connect(os.environ["DEMO_WAREHOUSE"]) as con: # one transaction
con.execute("CREATE TABLE IF NOT EXISTS orders_o (ds TEXT, order_id TEXT, amount REAL)")
con.execute("DELETE FROM orders_o WHERE ds = ?", (ds,)) # replace this run's slice only
con.executemany("INSERT INTO orders_o VALUES (?, ?, ?)", source_rows(ds))
load()
orders_append()
orders_overwrite()
''')
bag = DagBag(dag_folder=DAGS_FOLDER)
march_1 = pendulum.datetime(2026, 3, 1, tz="UTC")
for dag_id in ("orders_append", "orders_overwrite"):
for attempt in (1, 2):
run_dag(bag.get_dag(dag_id), logical_date=march_1)
with sqlite3.connect(WAREHOUSE) as con:
for table in ("orders_a", "orders_o"):
rows, total = con.execute(f"SELECT COUNT(*), SUM(amount) FROM {table} WHERE ds = '2026-03-01'").fetchone()
print(f"{table}: {rows} rows, total {total}")
orders_a: 6 rows, total 129.0
orders_o: 3 rows, total 64.5
The append version doubled the day’s revenue on the second run; the overwrite version still has the correct three rows. The same applies to files: write dt=2026-03-01/ as a unit (overwrite the folder or write elsewhere and swap) instead of adding files to it.
Pitfalls
- Using
datetime.now()or “the latest file” instead of the run’s dates, so a rerun processes a different slice. - Delete and insert in separate transactions: a failure in between leaves the partition empty.
- Side effects that cannot be undone, such as sending emails or calling a payment API. Make them conditional on a stored record of what was already done (an idempotency key), or move them to the very end.
- Overwriting more than the run’s slice, for example truncating the whole table in a daily task, which breaks backfills that run several days in parallel.
MERGEwithout a deterministic key, so reruns insert near-duplicates.
In interviews
“What does idempotent mean for a pipeline, and how do you make an Airflow task idempotent?” Define it (same final state however many times it runs for an interval), explain why Airflow needs it (retries, clears, backfills), and give the three patterns. A story where an append caused duplicates after a retry, and the fix, makes the answer stand out.
Avoiding top-level code
What it is
Top-level code is everything in a DAG file outside task callables: imports, constants, the DAG definition, and anything else that runs when the file is imported. The DAG processor imports each file repeatedly (by default it may reparse a file every 30 seconds, [dag_processor] min_file_process_interval), and each task start imports it again on the worker. Anything expensive at the top level is therefore paid over and over.
How it works
Keep the top level to imports, constants and DAG construction. Move everything else:
| Instead of this at the top level | Do this |
|---|---|
| API call or database query to decide which tasks exist | read a static config file generated in CI; or use dynamic task mapping at run time |
Variable.get("x") |
{{ var.value.x }} in a templated field, or Variable.get inside the task |
import pandas, import tensorflow |
import inside the task function |
| Creating hooks or clients | create them inside execute or the task |
| Reading large files | read them in a task |
The DagBag records how long each file took to import, which is a quick way to find offenders:
write_dag("slow_top_level.py", '''
import time
from datetime import datetime
from airflow.sdk import DAG
from airflow.providers.standard.operators.empty import EmptyOperator
time.sleep(1.5) # stands in for an API call or a heavy import at the top level
with DAG("slow_dag", schedule=None, start_date=datetime(2026, 1, 1), tags=["demo"],
default_args={"owner": "platform", "retries": 1}):
EmptyOperator(task_id="only_task")
''')
bag = DagBag(dag_folder=DAGS_FOLDER)
for stat in sorted(bag.dagbag_stats, key=lambda s: s.duration, reverse=True):
print(f"{stat.file:22} {stat.duration.total_seconds():5.2f}s dags={stat.dags}")
slow_top_level.py 1.50s dags=['slow_dag']
orders_load.py 0.00s dags=['orders_append', 'orders_overwrite']
On a real deployment the same figures appear as the dag_processing.last_duration metric and with airflow dags report. A file that exceeds [core] dagbag_import_timeout (30 seconds by default) fails to import and all its DAGs disappear from the UI.
Pitfalls
- “It is only one small API call”: multiplied by every parse and every task start, across every DAG file.
- Helper modules imported by DAG files that do work at import time. The rule applies to everything the DAG file imports.
- Code that behaves differently at parse time and run time (for example reading the current date at the top level to build task ids).
- Generating hundreds of DAGs in one file from a slow source.
In interviews
“Why is my scheduler slow?” or “Why do my DAGs keep disappearing?” often comes back to top-level code. Explain the parsing loop, what to move where (the table above), and how you would find slow files (DagBag stats, airflow dags report, parse-time metrics).
Testing DAGs
What it is
A DAG test suite has three layers:
- Integrity tests: does every file import, without cycles, and does every DAG follow team rules (owner, tags, retries, no
catchupsurprises)? Cheap, run on every commit. - Unit tests for task logic and custom operators, hooks and timetables, without running Airflow.
- Integration tests: run a whole DAG (or a few tasks) with
dag.test()against test data or containers.
How it works
Integrity tests load the DAG folder into a DagBag, the same parser the DAG processor uses, and assert on import_errors and on each DAG’s attributes. Below are pytest-style tests; with pytest installed you would put them in tests/test_dag_integrity.py and run pytest tests/. Here they are called directly. First, add two broken files of the kind integrity tests exist to catch:
write_dag("cycle.py", '''
from datetime import datetime
from airflow.sdk import DAG
from airflow.providers.standard.operators.empty import EmptyOperator
with DAG("cycle_dag", schedule=None, start_date=datetime(2026, 1, 1)):
a, b = EmptyOperator(task_id="a"), EmptyOperator(task_id="b")
a >> b >> a
''')
write_dag("missing_dependency.py", '''
from datetime import datetime
import great_expectations_not_installed # a library the workers do not have
from airflow.sdk import DAG
''')
write_dag("no_owner.py", '''
from datetime import datetime
from airflow.sdk import DAG
from airflow.providers.standard.operators.empty import EmptyOperator
with DAG("no_owner_dag", schedule="@daily", start_date=datetime(2026, 1, 1)):
EmptyOperator(task_id="t")
''')
# tests/test_dag_integrity.py
DAG_BAG = DagBag(dag_folder=DAGS_FOLDER) # in a real suite: a session-scoped pytest fixture
def test_no_import_errors():
errors = {os.path.basename(f): e.strip().splitlines()[-1] for f, e in DAG_BAG.import_errors.items()}
assert not errors, f"import errors: {errors}"
def test_every_dag_has_owner_tags_and_retries():
bad = [dag_id for dag_id, d in DAG_BAG.dags.items()
if d.default_args.get("owner") in (None, "airflow")
or not d.tags
or d.default_args.get("retries", 0) < 1]
assert not bad, f"missing owner, tags or retries: {sorted(bad)}"
def test_files_parse_quickly():
slow = [s.file for s in DAG_BAG.dagbag_stats if s.duration.total_seconds() > 1.0]
assert not slow, f"slow to parse: {slow}"
def test_orders_overwrite_structure():
d = DAG_BAG.get_dag("orders_overwrite")
assert d.task_ids == ["load"]
assert d.catchup is False
for test in (test_no_import_errors, test_every_dag_has_owner_tags_and_retries,
test_files_parse_quickly, test_orders_overwrite_structure):
try:
test()
print(f"PASS {test.__name__}")
except AssertionError as exc:
print(f"FAIL {test.__name__}: {exc}")
FAIL test_no_import_errors: import errors: {'cycle.py': 'AirflowDagCycleException: Cycle detected in Dag: cycle_dag. Faulty task: b', 'missing_dependency.py': "ModuleNotFoundError: No module named 'great_expectations_not_installed'"}
FAIL test_every_dag_has_owner_tags_and_retries: missing owner, tags or retries: ['no_owner_dag']
FAIL test_files_parse_quickly: slow to parse: ['slow_top_level.py']
PASS test_orders_overwrite_structure
Each failure points to a real problem: a cycle and a missing library (both would show up as import errors on the DAG processor), a DAG without an owner, tags or retries, and a file that is slow to parse. The structure test shows the other kind of assertion: checking task ids and dependencies so that a refactor cannot silently drop a task.
Note that the DagBag only parses files that contain both the words airflow and dag, ignoring case (safe mode, [core] dag_discovery_safe_mode), so a helper module without them is not reported, and is not loaded as a DAG either.
Unit tests call task logic directly. For TaskFlow tasks use .function; for custom operators call execute() with a minimal context; keep business logic in plain functions so most tests need no Airflow at all:
from airflow.sdk import task
def clean_amounts(rows: list[dict]) -> list[dict]:
"""Plain function: easy to test, used by the task below."""
return [r for r in rows if r["amount"] is not None and r["amount"] >= 0]
@task
def clean(rows: list[dict]) -> list[dict]:
return clean_amounts(rows)
def test_clean_drops_negative_and_null():
rows = [{"amount": 10.0}, {"amount": -1.0}, {"amount": None}]
assert clean.function(rows) == [{"amount": 10.0}]
test_clean_drops_negative_and_null()
print("PASS test_clean_drops_negative_and_null")
PASS test_clean_drops_negative_and_null
Integration tests run the DAG with dag.test(), which executes every task in one process without a scheduler, and then assert on results, as the idempotency example did. In a DAG file, if __name__ == "__main__": dag.test() lets developers run python dags/my_dag.py; in CI, airflow dags test <dag_id> <date> does the same from the command line. Point connections at test systems with dag.test(conn_file_path=..., variable_file_path=...) or environment variables, and use containers (for example a PostgreSQL container) for the systems the DAG talks to.
Pitfalls
- Only testing that files import. Add rule checks and structure assertions; they catch the mistakes that pass import.
- Integrity tests that run against a different Airflow or provider version than production. Install the same versions (with the constraints file) in CI.
- Tests that need real cloud credentials, so they are skipped or flaky. Mock hooks, use containers or local stand-ins.
- Asserting on
datetime.now()-dependent output. Pass a fixedlogical_date. - Forgetting the two-run check for idempotency: run the DAG twice for the same date and compare.
In interviews
“How do you test Airflow DAGs?” Name the three layers, show a DagBag import-error test and a rule test (owner, tags, retries), say how you unit test task logic (.function, plain functions, operator execute with mocks), and how you run integration tests (dag.test(), airflow dags test, containers). Mentioning that dag.test() replaced the DebugExecutor in Airflow 3 is a nice touch.
CI/CD for DAGs
What it is
CI/CD for DAGs means every change goes through automatic checks and is deployed the same way every time: no editing files on the server, no “it worked on my laptop”. Because DAG files are code that the scheduler executes, a broken file in production can stop a whole team’s pipelines.
How it works
A typical pipeline:
- Lint and format:
ruff(itsAIRrules flag Airflow-specific problems, including imports and arguments removed or moved in Airflow 3), plus type checks if you use them. - Install the production versions: the same Airflow version, providers and libraries, pinned with the constraints file, often by building the same container image as production.
- Integrity and unit tests: the DagBag tests above and the unit tests, with
pytest. - Optional integration tests:
dag.test()for key DAGs against containers or a staging environment. - Deploy the DAG code. In Airflow 3, DAGs come from DAG bundles: a
LocalDagBundle(a folder, filled by your deployment, an image or a sync sidecar) or aGitDagBundle(from the Git provider), where Airflow itself fetches a branch or tag from the repository. Git bundles also give DAG versioning: runs record which bundle version they used, and you can choose whether reruns use the original or the latest version. - Post-deploy checks:
airflow dags list-import-errors(orGET /api/v2/importErrors) should be empty; alert if it is not.
# .github/workflows/dags.yml
name: dags
on:
pull_request:
push:
branches: [main]
jobs:
test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with:
python-version: "3.11"
- name: Install pinned Airflow, providers and test tools
run: |
pip install "apache-airflow==3.3.2" -r requirements.txt \
--constraint "https://raw.githubusercontent.com/apache/airflow/constraints-3.3.2/constraints-3.11.txt"
pip install pytest ruff
- name: Lint
run: ruff check dags/ plugins/ --select E,F,AIR
- name: Initialise a throwaway metadata DB
run: airflow db migrate
- name: DAG integrity and unit tests
run: pytest tests/ -q
deploy:
needs: test
if: github.ref == 'refs/heads/main'
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Publish DAGs
run: ./scripts/deploy_dags.sh # e.g. build and push the image, or tag the commit a GitDagBundle tracks
# airflow.cfg: read DAGs from Git instead of a local folder (requires the Git provider and a git connection)
[dag_processor]
dag_bundle_config_list = [
{"name": "team_dags",
"classpath": "airflow.providers.git.bundles.git.GitDagBundle",
"kwargs": {"git_conn_id": "github_dags", "tracking_ref": "main", "subdir": "dags", "refresh_interval": 300}}
]
Pitfalls
- Deploying DAG code and the libraries it needs separately, so for a while the DAG imports a package the workers do not have. Ship them together (one image), or make the DAG tolerate the old version.
- Different versions in CI and production, so tests pass and parsing fails.
- Deleting or renaming a
dag_idin a refactor; history is orphaned and dependants (assets, sensors) break. Treatdag_ids as an API. - Changing a DAG’s structure while runs are in flight. In Airflow 3 runs keep the DAG version they started with when bundles support versioning; with a plain local folder they do not.
- No check after deploy: an import error can sit unnoticed until the morning’s runs are missing.
In interviews
“How do you deploy DAGs safely?” Describe the pipeline: lint (including Airflow 3 migration rules), pinned installs, DagBag and unit tests, optional integration tests, deployment through an image or a Git DAG bundle, and a post-deploy import-error check. Mention DAG versioning in Airflow 3 and keeping dag_ids stable.
Practice questions
A daily task appends to a table, and after a retry the day’s numbers are doubled. How do you fix the task so this cannot happen?
Make it idempotent for its interval: in one transaction delete the rows for the run’s ds (or overwrite the partition) and insert them again, or MERGE on a unique key. Parameterise it by the run’s dates, never now(). Then a retry, clear or backfill produces the same final state.
Write a test that fails the build when any DAG file has an import error.
def test_no_import_errors_in_repo():
bag = DagBag(dag_folder=DAGS_FOLDER) # in a real repo: DagBag(dag_folder="dags/")
assert bag.import_errors, "this demo folder deliberately contains broken files"
# in your suite: assert not bag.import_errors, bag.import_errors
test_no_import_errors_in_repo()
print("ok")The DagBag uses the same parser as the DAG processor, so a file that fails here would fail in production. Run it in CI with the same Airflow and provider versions as production.
Your DAG file reads a list of tables from an API to create one task per table. What is wrong, and what are the alternatives?
The API is called on every parse (repeatedly per file) and every task start, which slows parsing, can hit the import timeout, and makes the DAG structure change without a code change. Alternatives: generate a static config file in CI and read it locally, or list the tables in a task at run time and use dynamic task mapping.
How do you unit test the logic inside a @task function?
Call the original function through .function (for example my_task.function(args)), or better, keep the logic in a plain function that the task calls and test that directly. No scheduler or database is needed.
What replaced the DebugExecutor for running a DAG in a debugger?
dag.test() (from a DAG file run with python dags/my_dag.py) or airflow dags test <dag_id> <date>. Both run all tasks in a single process, so breakpoints work. The DebugExecutor and SequentialExecutor were removed in Airflow 3.
What is a DAG bundle, and why does it matter for CI/CD in Airflow 3?
A DAG bundle is a source of DAG files that Airflow reads: a local folder, or a Git repository (GitDagBundle) that Airflow fetches itself. With a versioned bundle such as Git, each DAG run records the version it used, so in-flight runs are not changed by a deploy and reruns can choose the original or the latest code. Deployment can become “merge to the tracked branch”.
Key takeaways
- Idempotent tasks (partition overwrite, merge on key, write then swap) make retries, clears and backfills safe; test by running twice for the same date.
- Keep DAG files declarative: no API calls, queries or heavy imports at the top level, because files are parsed constantly.
- Test in three layers: DagBag integrity and rule tests, unit tests of task logic, and
dag.test()integration tests. - In CI, lint with Airflow-aware rules, install pinned production versions, run the tests, deploy through an image or a Git DAG bundle, and check for import errors afterwards.
- Treat
dag_ids and asset names as interfaces; changing them breaks history and dependants.
Progress is saved in this browser only. No account needed.

