skip to content

In Airflow, why does top-level code in a DAG file slow down the whole scheduler?

level: seniorimportance: should knowfreq 58%

answer

  1. module level runs far more often than you think
  2. the file is imported, not executed once
  3. every parse pays it again
  4. the cost lands on unrelated DAGs too
  5. push work down into the task body

basics

~20 s

Because everything at module level runs every time the file is parsed, and Airflow re-parses every DAG file on a short repeating interval. An API call or query at the top level is paid on every parse, for every file, delaying scheduling everywhere.

solid answer

~50 s

Airflow's parsing loop imports every file in the DAGs folder repeatedly, on the interval set by `min_file_process_interval`, to build and serialize the DAG structure into the metadata database. Anything at module level — a `requests.get` for config, a warehouse query listing tables, a heavy library import, a `Variable.get` — executes on **every** one of those parses, not once per run and not once per deployment. The cost multiplies across files and shows up as sluggish scheduling for unrelated DAGs; if a file exceeds `dagbag_import_timeout` it is dropped and its DAGs disappear from the UI. Non-deterministic top-level code is worse: `start_date=datetime.now()` makes the schedule shift on every parse. The fix is to move work into task bodies, which run on workers at run time, and use Jinja templates or `params` for values that must be resolved per run.

code

python · 10 lines
python
# BAD: executed on every parse of this file
import requests
from airflow.models import Variable

tables = requests.get("https://config.internal/tables").json()
api_key = Variable.get("vendor_api_key")

with DAG("ingest", schedule="@daily", start_date=START, catchup=False):
    for t in tables:
        PythonOperator(task_id=f"load_{t}", python_callable=load, op_args=[t, api_key])

go deeper

for a junior

Know that code outside a task runs whenever Airflow reads the file, so calls to APIs or databases belong inside the task function, not at the top of the module.

for a middle

Explain the parse loop and serialization: the file is re-imported on an interval, its structure written to the metadata database, and every top-level statement is paid on each pass.

for a senior

Diagnose it in production — parse-duration reports, import timeouts making DAGs vanish, cross-DAG scheduling lag — and refactor the offending file into run-time work.

for a principal

Set the standard: parse-time budgets, lint or CI checks against top-level network and Variable access, and a shared pattern for config-driven DAGs that keeps a live dependency out of the parsing path.

## Two clocks: parse time and run time A DAG file is a Python program whose job is to *describe* a graph. It is executed by the parsing component (the scheduler's DAG file processor in Airflow 2; a mandatory standalone DAG processor in Airflow 3) again and again, whether or not any DAG run is happening — and often by the webserver's and workers' own imports as well. Run time is different. Task bodies execute exactly once per task instance, on a worker, when the scheduler decides that instance is eligible. Almost every Airflow performance question reduces to knowing which of these two clocks a given line of code is on. ## The parsing loop The processor scans the DAGs folder for files (`dag_dir_list_interval` controls how often it looks for new ones), imports each one, collects the `DAG` objects in its globals, checks for cycles, and writes a JSON serialization into the metadata database. `min_file_process_interval` sets the minimum gap before the same file is parsed again — short enough that the loop feels continuous. The serialized form is what the scheduler and the UI read afterwards, which is why the UI can render a DAG without importing your code. So the arithmetic is brutal: a two-second top-level HTTP call, in one file, parsed continuously, is two seconds of processor time burnt repeatedly, forever. Multiply by a hundred files with the same habit and DAGs across the whole deployment start being scheduled late, with no single DAG obviously at fault. There is also an import timeout (`dagbag_import_timeout`); a file that exceeds it is abandoned and its DAGs vanish from the UI until a parse succeeds — an outage whose cause is invisible unless you look at import errors. ## What belongs where ```python # BAD - executes on every parse of this file import requests tables = requests.get("https://config.internal/tables").json() with DAG("ingest", schedule="@daily", start_date=START, catchup=False): for t in tables: PythonOperator(task_id=f"load_{t}", python_callable=load, op_args=[t]) ``` ```python # GOOD - the call happens on a worker, once per run @task def fetch_tables() -> list[str]: return requests.get("https://config.internal/tables").json() load.expand(table=fetch_tables()) ``` The general rule: at module level, only cheap, deterministic, local expressions — constants, small literal lists, a config file already on disk, and the DAG and operator constructors themselves. Everything that touches a network, a database, a secrets backend or a large library goes inside a task body or its callable, not at import. Two refinements. First, **imports count**: `import pandas` at the top of a DAG file is paid on every parse of that file, so move heavy imports inside the callable. Second, **connections and Variables**: `Variable.get("x")` at module level hits the metadata database on every parse; fetch it inside the task, or reference it through a Jinja template such as `{{ var.value.x }}`, which is rendered at run time rather than at parse time. ## Determinism matters as much as cost Top-level code that changes between parses gives you a DAG whose structure is unstable. `start_date=datetime.now()` is the canonical example: every parse moves the start date, so the scheduler's idea of which intervals exist keeps shifting. A task list derived from a live API means tasks silently appearing and disappearing between parses, leaving orphaned task instances from earlier runs. Pin the structure to something stable — a fixed timestamp, a config file committed to the repo — and push volatility into run time. ## Diagnosing it Airflow gives you direct measurements. `airflow dags report` shows per-file parse duration and how many DAGs each produced; the scheduler exposes DAG-processing metrics; and the import-errors view lists files that timed out or raised. Sort by parse duration and the offenders are usually obvious — one file taking seconds while the rest take milliseconds. Mitigations, in order: remove the work; if it genuinely must be parse-time, cache it into a file the parse reads cheaply; raise `min_file_process_interval` to reduce frequency; and split enormous DAG files so a slow one does not delay the rest. ## The interview answer Say the two-clocks framing first, then the loop, then the concrete consequences — scheduler lag across unrelated DAGs, import timeouts making DAGs disappear, structural instability — and finish with the fix. That progression shows you have operated Airflow rather than only authored DAGs on it.

  • A DAG file needs a list of tables from an external service to decide how many tasks exist. How do you handle it without top-level calls?
    Prefer run-time fan-out: an upstream task fetches the list and a mapped task expands over it, so the network call happens once per run on a worker. If the graph genuinely must be shaped at parse time, have a separate scheduled job write the list to a file or Variable and let the parse read that cheap, local source instead of calling the service.
  • How would you find which DAG file is slowing down parsing?
    Run `airflow dags report`, which lists per-file parse duration and DAG counts, and check DAG-processor metrics and the import-errors view for files that timed out. Offenders usually stand out by an order of magnitude. Confirm by importing the file yourself and timing it, then look for network calls, Variable lookups or heavy imports at module level.
  • Why is `start_date=datetime.now()` in a DAG definition a bug?
    It is evaluated on every parse, so the start date keeps moving forward and the set of intervals the scheduler believes exists is never stable — runs can be missed or behave inconsistently between parses. Use a fixed timestamp such as `pendulum.datetime(2024, 1, 1, tz="UTC")` so the schedule is deterministic across parses and deployments.

saying these in an interview costs you the question

  • Thinks a DAG file is imported once at deployment
  • Fetches Variables or connections at module level
  • Uses datetime.now() for start_date or in the DAG structure
  • Blames the executor when the symptom is slow parsing
  • Puts heavy library imports at the top of every DAG file

context