Airflow courseLesson 6 of 8
Airflow course · Lesson 6 of 8
Operating Airflow: Backfills, Retries, Alerting, Logging and the REST API
Run Airflow day to day: backfill history, clear tasks safely, set retries and alerts, configure logs and StatsD metrics, and use the REST API v2.
On this page
- Running the examples
- Backfills and reruns
- What it is
- How it works
- Pitfalls
- In interviews
- Clearing tasks
- What it is
- How it works
- Pitfalls
- In interviews
- Retries and alerting
- What it is
- How it works
- Pitfalls
- In interviews
- Logging configuration
- What it is
- How it works
- Pitfalls
- In interviews
- Metrics with StatsD
- What it is
- How it works
- Pitfalls
- In interviews
- Airflow REST API
- What it is
- How it works
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
Writing DAGs is half the job; the other half is running them: reprocessing history after a bug fix, rerunning a failed task, being told when something breaks, and finding out why. This lesson covers the operational tools in Airflow 3: backfills, clearing, retries and alerts, logs and metrics, and the REST API that lets other systems drive Airflow.
Running the examples
The examples use Airflow 3.3 with a metadata database created by airflow db migrate. run_dag runs a DAG once with dag.test() as in earlier lessons. For the REST API section the page also starts Airflow’s API application in-process with FastAPI’s TestClient, so real /api/v2 requests run without a server; the two environment variables below configure the built-in simple auth manager with one admin user.
import contextlib, io, json, os, tempfile
from datetime import datetime, timedelta
PASSWORDS_FILE = os.path.join(tempfile.mkdtemp(), "passwords.json")
with open(PASSWORDS_FILE, "w") as fh:
json.dump({"admin": "example-password"}, fh)
os.environ["AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_USERS"] = "admin:admin" # user:role
os.environ["AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_PASSWORDS_FILE"] = PASSWORDS_FILE
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
def run_dag(dag, **test_kwargs):
"""Register an in-memory DAG, run it once with dag.test(), return (run state, 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:
states = {ti.task_id: str(ti.state) for ti in dr.get_task_instances(session=session)}
return str(dr.state), dict(sorted(states.items()))
Backfills and reruns
What it is
A backfill creates DAG runs for a range of past logical dates on purpose: after adding a new DAG that should cover last quarter, after fixing a bug in a transformation, or after an outage. A rerun reprocesses runs that already exist, usually by clearing them.
How it works
Airflow 3 rebuilt backfills. They are now managed by the scheduler (Airflow 2’s airflow dags backfill ran the whole backfill inside the CLI process, which died if your terminal did), and you can create them from the UI, the CLI or the REST API.
# Preview which runs would be created
airflow backfill create --dag-id orders_daily \
--from-date 2026-03-01 --to-date 2026-03-31 --dry-run
# Create them: rerun dates whose runs failed, at most 4 runs at a time
airflow backfill create --dag-id orders_daily \
--from-date 2026-03-01 --to-date 2026-03-31 \
--reprocess-behavior failed --max-active-runs 4
| Option | Effect |
|---|---|
--from-date, --to-date |
the logical date range, both inclusive |
--reprocess-behavior none (default) |
only create runs for dates that have no run yet |
--reprocess-behavior failed |
also rerun dates whose existing run failed |
--reprocess-behavior completed |
rerun every date, including successful ones |
--max-active-runs |
concurrency limit for this backfill |
--run-backwards |
newest dates first (not allowed with depends_on_past) |
--dag-run-conf |
JSON conf for every run |
Backfill runs have the run type backfill, can be paused, unpaused and cancelled (/api/v2/backfills/{id}/pause, /unpause, /cancel), and respect the DAG’s normal concurrency limits and pools. For a single date you can also just trigger a run with that logical date.
Pitfalls
- Non-idempotent tasks. A backfill of a task that appends rows duplicates data. Overwrite or merge per interval.
- Backfilling with
--reprocess-behavior completedwhen you meantfailed, and rewriting a month of good data. - Hammering a source system: set
--max-active-runsand use pools. depends_on_past=Truemakes a backfill strictly sequential; one failure stops everything after it.- Backfilling a DAG whose code changed: by default the runs use the latest DAG version, so make sure today’s code can process old data (schemas, columns) correctly.
In interviews
“You found a bug that corrupted the last two weeks. How do you fix the data?” Fix and deploy the code, make sure the tasks are idempotent, run a dry-run backfill for the range with reprocess-behavior completed, limit concurrency, then validate. Mention that Airflow 3 backfills run in the scheduler and can be managed from the UI and API, unlike the Airflow 2 CLI-process backfill.
Clearing tasks
What it is
Clearing a task instance resets its state so the scheduler runs it again, in its existing DAG run. It is how you rerun one failed task, a task and everything downstream of it, or a whole run, without creating new runs.
How it works
When you clear from the UI, the CLI (airflow tasks clear) or the API, you choose the scope:
| Option | CLI flag | Clears also |
|---|---|---|
| Downstream | -d |
every task downstream of the selected ones (usually what you want) |
| Upstream | -u |
every task upstream |
| Only failed | -f |
restricts to failed tasks |
| Past / Future | API include_past / include_future |
the same tasks in earlier or later runs |
| Date range | -s, -e |
runs with logical dates in the range |
Clearing resets the try count for the new attempts, removes the task’s XComs, and sets a finished DAG run back to queued (or running) so the scheduler picks it up. Clearing an ExternalTaskMarker can clear dependent tasks in other DAGs. Mark success and mark failed are the opposite tool: they change the state without running anything, for example to unblock downstream tasks after you fixed data by hand.
# Rerun the failed "load" task and everything after it for March, without prompting
airflow tasks clear orders_daily -t "^load$" -d -f -s 2026-03-01 -e 2026-03-31 -y
The REST API offers a dry run that shows exactly what would be cleared; there is an example in the REST API section below.
Pitfalls
- Clearing without “downstream”: the fixed task reruns but downstream tasks keep their old (possibly wrong) results.
- Clearing a sensor-heavy DAG at once, flooding slots.
-tis a regular expression;loadalso matchesload_customers. Anchor it (^load$).- Clearing a running task kills it (the task goes to
restartingand runs again), which may leave external jobs running; implementon_killin custom operators. - Using mark success to hide a real failure; downstream data may now be wrong.
In interviews
“How do you rerun just one failed task?” Clear it with downstream selected (UI, airflow tasks clear -t ... -d, or the API). Explain what clearing resets (state, tries, XComs) and when you would mark success instead. Always say the task must be idempotent.
Retries and alerting
What it is
Retries re-execute a failed task automatically, which absorbs transient failures such as timeouts and throttling. Alerting tells people when retries are not enough. In Airflow you implement it with callbacks, notifiers, and (for timing) deadline alerts.
How it works
Retry settings, usually in default_args:
| Argument | Meaning |
|---|---|
retries |
extra attempts after the first ([core] default_task_retries, 0 by default) |
retry_delay |
wait before the next try (5 minutes by default) |
retry_exponential_backoff |
in Airflow 3 a multiplier: 0 disables it, 2.0 doubles the delay each try (Airflow 2 took True/False) |
max_retry_delay |
cap on the delay when backing off |
execution_timeout |
fail a try that runs longer than this (then retry if allowed) |
Raise AirflowFailException for errors that retrying cannot fix (bad input, permission denied) so the task fails at once, and AirflowSkipException to skip.
Callbacks run when a task or DAG run changes state: on_failure_callback, on_retry_callback, on_success_callback, on_execute_callback and on_skipped_callback on tasks; on_failure_callback and on_success_callback on the DAG. Each receives the context (task instance, DAG run, exception). Notifiers are ready-made callbacks from providers, such as SlackNotifier (Slack provider) or SmtpNotifier (SMTP provider); pass an instance instead of a function.
from airflow.sdk import dag, task
events = []
def task_failed(context):
events.append(("task failed", context["ti"].task_id, str(context["exception"])))
def task_retrying(context):
events.append(("retrying", context["ti"].task_id, f"after try {context['ti'].try_number}"))
def task_succeeded(context):
events.append(("task succeeded", context["ti"].task_id))
def run_failed(context):
events.append(("dag run failed", context["dag_run"].dag_id, context.get("reason")))
@dag(
dag_id="orders_ops_demo",
schedule="@daily",
start_date=datetime(2026, 3, 1),
catchup=False,
on_failure_callback=run_failed, # DAG-level
default_args={
"retries": 1,
"retry_delay": timedelta(seconds=1), # minutes in real life
"retry_exponential_backoff": 2.0,
"max_retry_delay": timedelta(minutes=30),
},
)
def orders_ops_demo():
alerts = {"on_failure_callback": task_failed, # task-level callbacks
"on_retry_callback": task_retrying,
"on_success_callback": task_succeeded}
@task(**alerts)
def extract() -> str:
return "s3://lake/orders/dt=2026-03-01/"
@task(**alerts)
def load(path: str) -> None:
raise RuntimeError("warehouse unavailable")
load(extract())
ops_dag = orders_ops_demo()
print(run_dag(ops_dag, logical_date=pendulum.datetime(2026, 3, 1, tz="UTC")))
for event in events:
print(event)
('failed', {'extract': 'success', 'load': 'failed'})
('task succeeded', 'extract')
('retrying', 'load', 'after try 1')
('task failed', 'load', 'warehouse unavailable')
('dag run failed', 'orders_ops_demo', 'task_failure')
Task callbacks can also go in default_args so every task gets them; here they are passed to each task explicitly. load failed, the retry callback fired, the second attempt failed too, and then both the task-level and DAG-level failure callbacks ran. A notifier is used the same way:
from airflow.providers.slack.notifications.slack import SlackNotifier
on_failure_callback=SlackNotifier(
slack_conn_id="slack_alerts",
text="{{ dag.dag_id }}.{{ ti.task_id }} failed for {{ ds }}: {{ ti.log_url }}",
channel="#data-alerts",
)
Pitfalls
- Retrying non-idempotent tasks (duplicate rows) or permanent errors (wasted time). Classify errors.
- No
execution_timeout: a hung task never fails, so no retry and no alert happens. - Alerting on every retry instead of the final failure: alert fatigue. Alert on failure, log on retry.
- Callbacks that raise exceptions or are slow. They run inside Airflow components; keep them short and wrapped in try/except.
- Email alerts without SMTP configured.
email_on_failureneeds the SMTP provider and[smtp]settings; most teams use chat or paging notifiers instead. - Expecting callbacks for “this should have finished by now”. That is a deadline alert (see the scheduling lesson).
In interviews
“How do you make a pipeline resilient and make sure people hear about failures?” Retries with backoff and a cap, execution_timeout, AirflowFailException for permanent errors, idempotent tasks, failure callbacks or notifiers to the on-call channel, and deadline alerts for late runs. Mention the Airflow 3 change to retry_exponential_backoff if you have used it.
Logging configuration
What it is
Every task try writes a task log, which the UI shows next to the task instance. Components (scheduler, DAG processor, triggerer, API server) write their own logs. Configuring logging is mostly about where task logs are stored and how long they live.
How it works
- Local logs go under
[logging] base_log_folder($AIRFLOW_HOME/logsby default), with one file per try. The default[logging] log_filename_templatein 3.3 isdag_id=<dag>/run_id=<run>/task_id=<task>/[map_index=<n>/]attempt=<try>.log. - With distributed executors workers are often short-lived (always, with the KubernetesExecutor), so turn on remote logging: logs are uploaded to object storage when the task finishes and read back by the UI.
- Inside a task, use the task logger (
self.login operators,logging.getLogger(__name__)orprintin TaskFlow, which is captured). Airflow masks values it knows to be secrets (connection passwords, sensitive variables) in task logs. - Log levels:
[logging] logging_level(INFOby default). A customlogging_config_classlets you change handlers and formats.
[logging]
remote_logging = True
remote_base_log_folder = s3://my-airflow-logs/prod # or gs://..., wasb://...
remote_log_conn_id = aws_logs # a connection with write access
delete_local_logs = True # keep worker disks clean
logging_level = INFO
Elasticsearch and OpenSearch handlers are available through providers for teams that ship logs to a search cluster; the UI then links to or reads from it.
Pitfalls
- KubernetesExecutor without remote logging: the logs disappear with the pod.
- Logs on local disks filling up workers; enable
delete_local_logswith remote logging, or rotate. - Printing secrets that Airflow does not know about (for example a token you built yourself); register them with
mask_secret()fromairflow.sdk.log, or simply never log them. - Debug-level logging everywhere in production: large volumes and slower tasks.
In interviews
Expect “Where do Airflow logs go, and what changes on Kubernetes?” Local per-try files by default, remote logging to object storage for distributed or ephemeral workers, the connection and path settings, and secret masking.
Metrics with StatsD
What it is
Airflow emits operational metrics (counters, gauges and timers) about its own behaviour: scheduler loop duration, task successes and failures, pool usage, executor slots, DAG parsing times. They are sent over StatsD (or OpenTelemetry) to your monitoring system, for example a StatsD exporter scraped by Prometheus, or Datadog.
How it works
# pip install 'apache-airflow[statsd]'
[metrics]
statsd_on = True
statsd_host = statsd-exporter.monitoring # default localhost
statsd_port = 8125 # default 8125
statsd_prefix = airflow # default airflow
metrics_allow_list = scheduler,executor,pool,dagrun,ti,dag_processing # optional prefix filter
Set [metrics] otel_on = True (with the OpenTelemetry settings) instead to export through OTLP; Datadog tags are available with statsd_datadog_enabled.
Useful metrics to dashboard and alert on (names as in Airflow 3.3, prefixed with airflow.):
| Metric | Type | Why it matters |
|---|---|---|
scheduler_heartbeat |
counter | a flat line means the scheduler is down |
scheduler.scheduler_loop_duration |
timer | a growing loop means scheduling is falling behind |
scheduler.tasks.starving |
gauge | tasks that cannot run because of pool limits |
executor.open_slots, executor.queued_tasks, executor.running_tasks |
gauge | executor capacity and backlog |
pool.open_slots.<pool>, pool.queued_slots.<pool> |
gauge | which pools are saturated |
ti_failures, ti_successes, ti.finish.<dag>.<task>.<state> |
counter | failure rates |
dagrun.duration.success.<dag>, dagrun.duration.failed.<dag> |
timer | run time trends |
dagrun.schedule_delay.<dag> |
timer | how late runs start compared with their schedule |
dag_processing.import_errors, dag_processing.total_parse_time |
gauge | broken DAG files and slow parsing |
triggerer.capacity_left.<host>, triggers.blocked_main_thread |
gauge / counter | triggerer saturation and blocking triggers |
Pitfalls
- Turning on StatsD without installing the
statsdextra. - High-cardinality metric names (per DAG and task) overwhelming the monitoring system; use the allow and block lists.
- Watching only task failures. A dead scheduler produces no failures at all, so alert on
scheduler_heartbeatand on schedule delay.
In interviews
“How do you monitor Airflow itself?” StatsD or OpenTelemetry metrics to Prometheus or Datadog, with alerts on scheduler heartbeat, loop duration, queued and starving tasks, pool saturation, import errors and schedule delay; plus the /api/v2/monitor/health endpoint for liveness checks. Separating “Airflow is healthy” from “my pipelines are healthy” is a strong point.
Airflow REST API
What it is
The REST API lets other systems and scripts do what the UI does: trigger runs, read states and logs, clear tasks, manage connections, variables and pools, and create backfills. In Airflow 3 the stable public API lives under /api/v2, served by the airflow api-server process. Airflow 2’s /api/v1 was removed, along with the old experimental API that Airflow 2 already disabled by default.
How it works
- Authentication is token-based: POST credentials to the auth manager’s token endpoint (
/auth/tokenfor the simple and FAB auth managers) to get a JWT access token, then sendAuthorization: Bearer <token>. Which login methods exist depends on the configured auth manager. - Main resources:
/api/v2/dags,/dags/{dag_id}/dagRuns,/dagRuns/{run_id}/taskInstances,/clearTaskInstances,/backfills,/connections,/variables,/pools,/importErrors,/monitor/health,/version. - For scripts there is an official Python client package (
apache-airflow-client) generated from the same API specification.
from fastapi.testclient import TestClient
from airflow.api_fastapi.app import create_app
api = TestClient(create_app()) # the same app `airflow api-server` serves
token = api.post("/auth/token", json={"username": "admin", "password": "example-password"}).json()["access_token"]
auth = {"Authorization": f"Bearer {token}"}
print("health:", api.get("/api/v2/monitor/health").json()["metadatabase"])
print("v1:", api.get("/api/v1/dags").status_code, api.get("/api/v1/dags").json()["error"])
dag_info = api.get("/api/v2/dags/orders_ops_demo", headers=auth).json()
print("dag:", dag_info["dag_id"], "paused:", dag_info["is_paused"])
# Inspect the failed run created by run_dag above
runs = api.get("/api/v2/dags/orders_ops_demo/dagRuns", headers=auth,
params={"logical_date_gte": "2026-03-01T00:00:00Z", "logical_date_lte": "2026-03-01T00:00:00Z"}).json()
failed_run = runs["dag_runs"][0]
print("run:", failed_run["logical_date"], failed_run["state"])
# What would "clear failed tasks in that run" rerun? (dry run, nothing changes)
preview = api.post("/api/v2/dags/orders_ops_demo/clearTaskInstances", headers=auth,
json={"dry_run": True, "only_failed": True, "dag_run_id": failed_run["dag_run_id"]}).json()
print("would clear:", [(ti["task_id"], ti["state"]) for ti in preview["task_instances"]])
# Trigger a manual run with conf; Airflow 3 allows logical_date to be null
new_run = api.post("/api/v2/dags/orders_ops_demo/dagRuns", headers=auth,
json={"logical_date": None, "conf": {"mode": "full"}}).json()
print("triggered:", new_run["run_type"], new_run["state"], new_run["logical_date"], new_run["conf"])
# Which runs would a backfill of 1-4 March create, rerunning failed dates?
plan = api.post("/api/v2/backfills/dry_run", headers=auth, json={
"dag_id": "orders_ops_demo", "from_date": "2026-03-01T00:00:00Z",
"to_date": "2026-03-04T00:00:00Z", "reprocess_behavior": "failed"}).json()
print("backfill would create:", [b["logical_date"][:10] for b in plan["backfills"]])
health: {'status': 'healthy'}
v1: 404 /api/v1 has been removed in Airflow 3, please use its upgraded version /api/v2 instead.
dag: orders_ops_demo paused: True
run: 2026-03-01T00:00:00Z failed
would clear: [('load', 'failed')]
triggered: manual queued None {'mode': 'full'}
backfill would create: ['2026-03-01', '2026-03-02', '2026-03-03', '2026-03-04']
A few things to notice: new DAGs start paused ([core] dags_are_paused_at_creation is True), so the triggered run stays queued until the DAG is unpaused (PATCH /api/v2/dags/{dag_id} with {"is_paused": false}); the run triggered with a null logical date has none; and the backfill plan includes 1 March because its existing run failed and the reprocess behaviour is failed. The same calls over HTTP against a running API server:
TOKEN=$(curl -s -X POST https://airflow.example.com/auth/token \
-H 'Content-Type: application/json' \
-d '{"username": "svc_ci", "password": "'"$AIRFLOW_PASSWORD"'"}' | jq -r .access_token)
curl -s -X POST https://airflow.example.com/api/v2/dags/orders_daily/dagRuns \
-H "Authorization: Bearer $TOKEN" -H 'Content-Type: application/json' \
-d '{"logical_date": "2026-03-01T00:00:00Z", "conf": {"mode": "full"}}'
curl -s -X POST https://airflow.example.com/api/v2/backfills \
-H "Authorization: Bearer $TOKEN" -H 'Content-Type: application/json' \
-d '{"dag_id": "orders_daily", "from_date": "2026-03-01T00:00:00Z", "to_date": "2026-03-31T00:00:00Z", "reprocess_behavior": "failed", "max_active_runs": 4}'
Pitfalls
- Old integrations calling
/api/v1or basic-auth endpoints after an upgrade; they get 404 and must move to/api/v2with tokens. - Triggering by API with a logical date that already has a run: the request is rejected as a conflict. Omit or null the logical date for ad hoc runs.
- Long-lived personal tokens in CI. Use a service account and short-lived tokens.
- Polling run state in a tight loop; poll with backoff, or use asset events or callbacks.
- Field names changed between v1 and v2 (for example
execution_datebecamelogical_date, andis_pausedupdates usePATCH).
In interviews
“How would an external system trigger a DAG?” POST /api/v2/dags/{dag_id}/dagRuns with a bearer token and conf, or update an asset so asset-scheduled DAGs run. Mention authentication through the auth manager’s token endpoint, that /api/v1 is gone in Airflow 3, and the health endpoint for monitoring.
Practice questions
How do backfills differ between Airflow 2 and Airflow 3?
In Airflow 2, airflow dags backfill ran the backfill inside the CLI process, outside the scheduler, so it stopped if the process died and was not visible as a managed object. In Airflow 3, airflow backfill create (or the UI or POST /api/v2/backfills) creates a backfill that the scheduler runs, with a reprocess behaviour (none, failed, completed), its own max active runs, dry runs, and pause, unpause and cancel operations.
A task in yesterday’s run failed after the data it reads was fixed. How do you rerun it and its dependants only?
Clear that task instance with “downstream” (and “only failed” if appropriate) from the UI, with airflow tasks clear dag_id -t '^task$' -d -s <date> -e <date>, or with POST /api/v2/dags/{dag_id}/clearTaskInstances. The DAG run goes back to running and the scheduler reruns the cleared tasks. The tasks must be idempotent.
What does retry_exponential_backoff=2.0 with retry_delay=5 minutes and max_retry_delay=30 minutes do?
The wait grows by a factor of 2 per retry: about 5, 10, 20 minutes, then capped at 30 minutes. In Airflow 3 the argument is a multiplier (0 disables backoff); in Airflow 2 it was a boolean.
Your on-call channel is flooded with alerts for retries that later succeed. What do you change?
Send notifications from on_failure_callback (final failure) rather than on_retry_callback, log retries instead, and add deadline alerts for runs that are late. Make sure permanent errors raise AirflowFailException so they fail at once instead of retrying.
Why do task logs disappear with the KubernetesExecutor, and how do you fix it?
Each task runs in a pod that is deleted when it finishes, taking its local log file with it. Enable remote logging (remote_logging = True, remote_base_log_folder pointing at object storage, and remote_log_conn_id) so logs are uploaded when the task ends and read from storage by the UI.
Which metrics would you alert on to know that Airflow itself, not a pipeline, is unhealthy?
Scheduler heartbeat stopping, scheduler loop duration growing, executor queued tasks and open slots, tasks starving in pools, DAG processing import errors and total parse time, triggerer capacity left, and schedule delay. Also probe /api/v2/monitor/health.
An integration that worked on Airflow 2 now gets 404 from /api/v1/dags/x/dagRuns. What changed?
Airflow 3 removed /api/v1; the stable API is /api/v2, authenticated with bearer tokens from the auth manager (for example POST /auth/token). Field names also changed, notably execution_date became logical_date.
Key takeaways
- Airflow 3 backfills are scheduler-managed objects with reprocess behaviours and dry runs; only backfill idempotent tasks.
- Clearing reruns tasks in existing runs; include downstream tasks and anchor task regexes.
- Retries absorb transient errors;
retry_exponential_backoffis a multiplier in Airflow 3; alert on final failure with callbacks or notifiers. - Use remote logging for distributed or ephemeral workers, and keep secrets out of logs.
- Export StatsD or OpenTelemetry metrics and alert on scheduler heartbeat, loop duration, queued and starving tasks.
- The REST API is
/api/v2with bearer tokens;/api/v1was removed in Airflow 3.
Progress is saved in this browser only. No account needed.

