Airflow interview questions in 2026 are about Airflow 3: assets instead of datasets, the Task SDK and Task Execution API, DAG versioning, scheduler-managed backfills and deadline alerts in place of SLAs. Interviewers still check the basics (DAGs, operators, data intervals, retries), but the questions that separate candidates are operational. Why is this DAG stuck in queued? Why does the scheduler fall behind? How do you backfill safely after an outage? How do you move a large Airflow 2 estate to Airflow 3 without breaking production? This guide collects 50 high-value questions with model answers, from fundamentals to ten production scenarios.
How to use this guide
Answers use current Airflow 3.x terminology and say "previously X" where something was renamed or removed. General orchestration theory (idempotency, Airflow versus Dagster versus Prefect, safe reprocessing of years of history) is covered in our data engineering interview questions. This page stays on what is specific to Airflow.
- Freshers and early-career engineers: components, DAG and task concepts, TaskFlow, data intervals, catchup, connections and XComs (Q1βQ10).
- Mid-level data engineers: Airflow 3 changes, dynamic task mapping, deferrable operators, pools, retries, testing and CI/CD (Q11βQ36).
- Senior and platform engineers: executors, scaling, observability, multi-team platforms, ML and RAG pipelines, and the scenario section (Q29βQ50).
Airflow ships minor releases often and some newer features are marked experimental. In an interview, say "check the release notes for the current status" rather than guessing. The code snippets below use the airflow.sdk authoring interface and were import-tested against an Airflow 3.3 installation.
- Airflow fundamentals (Q1βQ10)
- What changed in Airflow 3 (Q11βQ17)
- DAG authoring (Q18βQ28)
- Executors, security and operations (Q29βQ36)
- Airflow for ML and AI pipelines (Q37βQ40)
- Real-world scenario questions (Q41βQ50)
- Key takeaways
- Interview preparation checklist
- FAQ
Airflow fundamentals
1. What is Apache Airflow, and what is it not?
Answer: Airflow is an open-source platform for authoring, scheduling and monitoring batch workflows as code. You define a DAG (directed acyclic graph) of tasks in Python, and Airflow decides when each task should run, runs it through an executor, retries it on failure and records every attempt. It is an orchestrator: it coordinates work that usually happens elsewhere, such as a Spark job, a warehouse query, a dbt run or a Kubernetes pod.
It is not a data processing engine and not a streaming system. Loading a large dataset into a worker's memory, or using Airflow to process individual events with sub-second latency, are both misuses. Airflow 3 does add event-driven scheduling, but that means "start a DAG run when an event arrives", not "process a stream".
Interview tip: Saying "Airflow should tell other systems what to do, not do the heavy lifting itself" early in the interview frames every later answer well.
2. What are the main components of an Airflow 3 deployment?
Answer: A typical Airflow 3 deployment has six parts:
- Scheduler: decides which DAG runs and task instances should exist and hands runnable tasks to the executor. The executor logic runs inside the scheduler process.
- DAG processor: parses DAG files and stores serialized DAGs. In Airflow 3 it is a separate process you must run (
airflow dag-processor), even locally. - API server: replaces the Airflow 2 webserver (
airflow api-server). It serves the React UI, the REST API v2 and the Task Execution API that workers use. - Workers: run the task code, through whatever executor you chose (Local, Celery, Kubernetes, Edge and others).
- Triggerer: runs asynchronous triggers for deferred tasks, so waiting does not consume worker slots.
- Metadata database: PostgreSQL or MySQL in production, holding DAG runs, task instances, XComs, connections and history.
The big architectural change is that workers no longer talk to the metadata database directly. They talk to the API server.
3. Explain DAG, task, operator, sensor and hook.
Answer: A DAG is the workflow: tasks plus dependencies plus a schedule. A task is one node in that graph, and a task instance is one execution of a task for one DAG run. An operator is a reusable template for a task (for example BashOperator, SQLExecuteQueryOperator or a cloud provider's job operator). A sensor is a special operator that waits for a condition, such as a file landing or a partition appearing. A hook is a lower-level client that wraps a connection to an external system; operators use hooks internally, and you can use hooks directly inside a @task function.
Real-world example: An S3-to-warehouse load might use a sensor to wait for the file, an operator to run the COPY command, and a hook inside a small Python task to fetch row counts for a data quality check.
4. What is the TaskFlow API, and when would you still use classic operators?
Answer: TaskFlow lets you write tasks as decorated Python functions (@dag, @task). Return values are passed downstream automatically through XCom, and dependencies are inferred from function calls, which makes Python-heavy DAGs shorter and easier to test. Classic operators are still the right choice when a provider already ships a well-tested operator for the job, especially deferrable ones for long-running external work (a cloud batch job, a warehouse query, a Kubernetes pod). Mixing the two is normal: TaskFlow for glue logic, provider operators for the external work.
Interview tip: Mention that in Airflow 3 you import dag and task from airflow.sdk, not from airflow.decorators. The old paths are deprecated.
5. Walk through the life cycle of a task instance.
Answer: A task instance usually moves from no state to scheduled (dependencies met), to queued (handed to the executor), to running (a worker picked it up), and ends in success, failed or skipped. Along the way it can enter up_for_retry after a failure with retries left, up_for_reschedule for a sensor in reschedule mode, deferred while a trigger waits in the triggerer, or upstream_failed when a parent failed and the trigger rule does not allow it to run.
Knowing where a task is stuck tells you which component to look at. Stuck in scheduled points at concurrency limits or pools. Stuck in queued points at the executor or workers. Stuck in deferred points at the triggerer.
6. What are logical date and data interval, and what changed in Airflow 3?
Answer: For a scheduled run, the data interval is the period of data the run is responsible for, and the logical date identifies the run (previously called execution_date, which Airflow 3 removed). A daily run covering 1 March typically starts after 1 March ends, which confuses many beginners.
Airflow 3 changed three things interviewers like to ask about:
- Manually triggered runs and asset-triggered runs can have
logical_date=Noneand no data interval, so code that readsdata_interval_startmust handle its absence. - The
create_cron_data_intervalssetting now defaults to False, so a bare cron string usesCronTriggerTimetable(run at that time) instead of the older data-interval timetable. Teams that rely on interval semantics must set it explicitly before upgrading. - Multiple runs with no logical date are now allowed, which suits ML training, inference and ad hoc work.
7. What is catchup, and why did its default change?
Answer: With catchup enabled, when a DAG is unpaused Airflow creates runs for every missed interval between its start_date and now. In Airflow 3, catchup_by_default is False, so a DAG that does not set catchup explicitly only runs from the current interval. The change prevents the classic accident where someone deploys a DAG with an old start_date and floods the cluster with hundreds of historical runs.
Interview tip: Say that you always set catchup explicitly and use a backfill (Q43) when you actually want history, because a backfill gives you a date range, a reprocess policy and its own concurrency limit.
8. What are Connections and Variables, and how do you keep secrets out of DAG code?
Answer: A Connection stores how to reach an external system (host, login, password, extras) under a conn_id; hooks and operators look it up by ID. A Variable is a global key-value setting. Neither should be hardcoded in DAG files. In production, store them in a secrets backend such as AWS Secrets Manager, Azure Key Vault, Google Secret Manager or HashiCorp Vault, with the DAG referencing only the ID. Airflow masks values of sensitive fields in logs, but that is a safety net, not a design. Q31 covers the lookup order.
9. What are XComs, and what are their limits?
Answer: XComs ("cross-communications") let tasks pass small values to each other. Each XCom is keyed by DAG, run, task and key, and TaskFlow return values become XComs automatically. The default backend stores them in the metadata database, and the Airflow docs are explicit that they are for small amounts of data, not dataframes. Large XComs bloat the database, slow the UI and API, and make serialization a bottleneck.
Two Airflow 3 details matter. First, xcom_pull() without task_ids now pulls only from the current task, so old code that relied on "find this key anywhere in the run" breaks. Second, XComs are cleared when a task retries, so they cannot hold state across retries. For larger payloads, write the data to object storage and pass the path, or configure the object storage XCom backend from the common IO provider (Q44).
10. What are trigger rules, and when do you change the default?
Answer: A trigger rule decides when a task runs based on its upstream states. The default, all_success, runs a task only when every parent succeeded. You change it for specific patterns: all_done for a cleanup or notification task that must run regardless of outcome, one_failed for an alert path, none_failed_min_one_success to join after a branch where some paths are skipped, and all_done_min_one_success (added in 3.1) when at least one parent must have succeeded. For infrastructure cleanup, setup and teardown tasks (Q26) are usually clearer than custom trigger rules.
What changed in Airflow 3
11. What are the most important changes in Airflow 3?
Answer: Airflow 3.0 was released in April 2025 as the largest change since 2.0. The headline items:
| Area | Airflow 2 | Airflow 3 |
|---|---|---|
| Task execution | Workers access the metadata database directly | Task Execution API and Task SDK; workers talk to the API server |
| Data-aware scheduling | Datasets | Assets, the @asset decorator, asset aliases, later asset partitioning |
| DAG history | Only the latest code is visible | DAG versioning and DAG bundles |
| Backfills | A separate CLI process | Managed by the scheduler, available from the UI, API and CLI |
| UI and API | Flask UI, REST API v1 | React UI on FastAPI, REST API v2 |
| Missed-time alerts | SLAs | Deadline alerts (from 3.1) |
| Authoring imports | airflow.models, airflow.decorators | airflow.sdk |
Also removed: SubDAGs, the SequentialExecutor, the static hybrid executors, schedule_interval, execution_date and DAG/XCom pickling. Common operators like PythonOperator and BashOperator moved into the apache-airflow-providers-standard package.
12. What are the Task Execution API and the Task SDK, and why did Airflow restrict database access?
Answer: In Airflow 2 every worker connected to the metadata database and ran user code in the same process that could open database sessions. That caused two problems: too many database connections at scale, and a security problem, because any DAG author could read or change Airflow's own state. Airflow 3 introduces the Task Execution API on the API server. Workers report state, heartbeats and XComs and fetch connections and variables through that API, and the Task SDK (apache-airflow-task-sdk, imported as airflow.sdk) is the stable, lightweight runtime and authoring interface that task code uses.
Practical consequences: custom operators that queried the metadata database must move to the REST API (the Airflow Python client) or the context the SDK exposes; workers no longer need database credentials; and tasks can run in remote or isolated environments. From 3.3 the same contract allows experimental Java and Go task implementations, while the DAG itself stays in Python.
Interview tip: The docs note that DAG author code can still run with database access in the DAG processor and triggerer. Mentioning this shows you read the security model, not just the release blog.
13. What are assets, and how does asset-aware scheduling work?
Answer: An asset (previously a dataset) is a logical piece of data identified by a name and a URI, such as a table or an object-store prefix. A producer task declares the asset in its outlets; when the task succeeds, Airflow records an asset event, and any DAG scheduled on that asset is triggered. This replaces fragile cross-DAG patterns like ExternalTaskSensor chains with explicit data dependencies, and the UI shows lineage across DAGs.
from airflow.sdk import Asset, dag, task
claims = Asset(name="raw_claims",
uri="s3://lake/raw/claims/")
@task(outlets=[claims])
def load_claims(): ...
@dag(schedule=claims)
def claims_features(): ...
Note that name is the first positional argument of Asset in Airflow 3, so use keywords. The @asset decorator goes further and creates a single-task DAG that materialises the asset. Airflow 3.2 added asset partitioning, so a downstream run can be triggered only for the partition that changed (for example one date), and 3.3 expanded the partition mappers. Asset events also carry an extra payload for annotations such as row counts.
14. What are DAG versioning and DAG bundles?
Answer: Airflow 3 records structural versions of each DAG in the metadata database, so the Graph, Grid and Code views show the DAG as it was for a given historical run. DAG bundles control where DAG code comes from: LocalDagBundle (the default, pointing at the DAGs folder), GitDagBundle, and object-store bundles for S3 and GCS. The important nuance is that only versioned bundles, such as the Git bundle, let a run keep using the code it started with. With local or object-store bundles, tasks always run the latest code on disk.
Airflow 3.3 added the rerun_with_latest_version setting, which controls whether clears, reruns and backfills use the original bundle version or the latest. Defaults differ: clears and reruns default to the original version, backfills default to the latest.
15. Which features were removed in Airflow 3, and what replaces them?
Answer:
- SubDAGs: use TaskGroups for visual grouping, and assets or dynamic task mapping for reuse and fan-out.
- SLAs: replaced by deadline alerts (Q16).
- SequentialExecutor: use LocalExecutor, which supports SQLite for local development.
- CeleryKubernetesExecutor and LocalKubernetesExecutor: use the multiple-executor configuration.
schedule_intervalandtimetablearguments: use the singlescheduleargument.execution_dateand context keys likeprev_ds,next_ds,yesterday_ds: uselogical_dateand the data interval.- REST API v1: use the FastAPI-based
/api/v2. - Registering operators, hooks and executors through plugins: import them as normal Python classes.
airflow webserver,airflow db init,airflow db upgrade: useairflow api-serverandairflow db migrate.
The default auth manager is now the Simple auth manager; teams that need Flask AppBuilder authentication install the FAB provider and configure it.
16. SLAs were removed. How do deadline alerts work?
Answer: A deadline alert has three parts: a reference (when the run was queued, its logical date, a fixed datetime, or the average runtime of recent successful runs), an interval (positive or negative), and a callback that runs if the DAG run has not finished by reference plus interval. Callbacks can be notifiers (Slack, email and others) or your own function, as an AsyncCallback run by the triggerer or, from 3.2, a SyncCallback run by the executor. A DAG can define several alerts.
from datetime import timedelta
from airflow.sdk import (AsyncCallback, DeadlineAlert,
DeadlineReference, dag)
@dag(deadline=DeadlineAlert(
reference=DeadlineReference.DAGRUN_QUEUED_AT,
interval=timedelta(hours=1),
callback=AsyncCallback(notify_on_call)))
def nightly_settlement(): ...
Here notify_on_call is an async function you define in a module the triggerer can import. Unlike old SLAs, which were task-level and widely seen as unreliable, deadlines are evaluated for the DAG run and fire while the run is still late, not after the fact.
Interview tip: The docs still mark deadline alerts as experimental. Say so, and say you would pin the Airflow version and test the alert path before relying on it for a regulatory deadline.
17. What has been added since Airflow 3.0?
Answer: Each minor release has added platform features:
- 3.1: human-in-the-loop operators (
HITLOperator,ApprovalOperatorand related, in the standard provider) that pause a task until someone responds in the UI or API; deadline alerts; UI translations; a React plugin system; and a streaming endpoint to wait for a DAG run's result, for inference-style use. - 3.2: asset partitioning, multi-team deployments (experimental) with per-team DAGs, connections, pools and executors, synchronous deadline callbacks, and editing XComs from the UI.
- 3.3: a task and asset state store that persists key-value state across retries and runs, pluggable retry policies, more partition mappers, experimental Java and Go task SDKs, and a Deadlines page in the UI.
Interview tip: You do not need every item. Picking the two that matter for your work (for example the state store for incremental loads and partitioning for date-partitioned lakes) and explaining why is a stronger answer.
DAG authoring
18. Write a small TaskFlow DAG with retries and dynamic fan-out.
Answer: A daily load that lists new files, loads each one in parallel and publishes an asset when done:
from datetime import datetime, timedelta
from airflow.sdk import Asset, dag, task
raw = Asset(name="raw_claims",
uri="s3://lake/raw/claims/")
@dag(schedule="@daily",
start_date=datetime(2026, 1, 1),
catchup=False,
default_args={"retries": 3,
"retry_delay": timedelta(minutes=5)})
def claims_ingest():
@task
def list_files() -> list[str]:
return ["claims/a.parquet", "claims/b.parquet"]
@task(max_active_tis_per_dag=4)
def load_file(key: str) -> int:
return 100 # rows loaded
@task(outlets=[raw])
def summarise(counts: list[int]) -> int:
return sum(counts)
summarise(load_file.expand(key=list_files()))
claims_ingest()
Points to explain: list_files returns object keys, not file contents; load_file.expand creates one mapped task instance per key at runtime; max_active_tis_per_dag caps how many load in parallel; the summary task receives the collected results and emits the asset event. The return values are small, so XCom is fine here.
19. How does dynamic task mapping work, and what are its limits?
Answer: expand() creates one task instance per element of a list or dict that is only known at runtime; partial() fixes the arguments that stay the same; expand_kwargs() maps over a list of argument dictionaries when each instance needs several different inputs. Classic operators support the same pattern, though arguments like task_id, pool and queue must go in partial(). You can map over the output of an upstream task and chain mapped tasks.
Limits: the [core] max_map_length setting caps the number of mapped instances (1024 by default), each mapped instance is a full task instance in the scheduler and database, and the upstream list itself travels through XCom. For tens of thousands of tiny items, batch them (map over chunks of a few hundred keys) rather than mapping each item.
Interview tip: Dynamic task mapping is the usual answer to "how did you replace SubDAGs or generate tasks in a loop at parse time".
20. What are deferrable operators, and why does the triggerer matter?
Answer: A deferrable operator starts external work, then calls defer() with a trigger and gives up its worker slot. The trigger is small async Python code that runs in the triggerer process, where thousands of triggers can wait together. When the trigger fires, the task resumes on a worker to finish. By default deferred tasks do not occupy pool slots either, though a pool can be configured to count them.
Requirements: at least one triggerer must be running, and trigger code must never block (no synchronous network or file calls in run()), because one blocking trigger stalls every trigger in that process. Many provider operators accept deferrable=True, and [operators] default_deferrable sets the default for operators that support both modes.
Real-world example: Consider a retailer whose nightly DAG submits dozens of warehouse transformations that each run for a long time. With synchronous operators every one of them holds a worker slot while doing nothing; deferred, they hold none, and the same workers can run other DAGs.
21. Compare sensor modes: poke, reschedule and deferrable.
Answer:
| Mode | Holds a worker slot while waiting? | Good for |
|---|---|---|
| poke | Yes, for the whole wait | Very short waits with frequent checks |
| reschedule | No, only during each check | Waits of minutes to hours with infrequent checks |
| deferrable | No; the triggerer waits | Long waits, many concurrent sensors, provider support available |
Always set timeout (and think about soft_fail), otherwise a sensor waiting for a file that never arrives can wait for days. With TaskFlow, @task.sensor(poke_interval=60, timeout=3600, mode="reschedule") turns a function returning a boolean into a sensor.
22. How do pools, priority weights and concurrency settings interact?
Answer: Several limits decide whether a runnable task actually starts:
[core] parallelism: maximum running task instances per scheduler.max_active_runsandmax_active_taskson the DAG: concurrent runs and concurrent tasks for that DAG.max_active_tis_per_dagon a task: concurrent instances of that task across runs.- Pools: named slot budgets, usually one per fragile downstream system. Tasks without a pool use
default_pool(128 slots by default). A heavy task can take several slots withpool_slots. - Priority weight and
weight_rule: order which queued tasks get free slots first.
In Airflow 3, priority weight is capped by available pool slots, so a high-priority task cannot starve everything else when a pool is contended. The usual design is a pool per protected resource (a core banking API, a licence-limited database), priorities for business-critical DAGs, and DAG-level limits to stop one DAG flooding the cluster.
23. How do you design retries properly?
Answer: Set retries, retry_delay and, for flaky external APIs, retry_exponential_backoff with a max_retry_delay. Add execution_timeout so a hung task fails instead of running forever. Most important, retries are only safe if the task is idempotent (Q25). Not every failure deserves a retry: an authentication error or bad input will fail the same way each time. Airflow 3.3 added pluggable retry policies, so a task can retry only on specific exception types or decide its backoff with custom logic, rather than retrying blindly.
Interview tip: Say that you separate transient errors (timeouts, throttling) from permanent ones and alert immediately on the second kind.
24. What is "top-level code", and why does it hurt performance?
Answer: Anything outside task callables runs every time the DAG processor parses the file, which happens regularly (by default each file at most every 30 seconds via min_file_process_interval). Database queries, API calls, Variable.get() calls or heavy imports at the top level slow parsing for every DAG, can hit parse timeouts, and hammer external systems. Move expensive work into tasks, import heavy libraries inside the task function, use Jinja templates such as {{ var.value.x }} (which resolve at run time), and keep DAG-generation logic deterministic.
25. How do you make an Airflow task idempotent?
Answer: A task is idempotent if running it twice for the same run gives the same result. Patterns: write to a partition keyed by the data interval and overwrite it (delete-then-insert in a transaction or INSERT OVERWRITE); use MERGE on a business key rather than plain inserts; derive dates from the run context, never from datetime.now(); make external calls with an idempotency key where the API supports one; and write files to deterministic paths. In Airflow 3, remember that manual and asset-triggered runs may have no data interval, so decide what such a run should process.
26. When do you use TaskGroups, and what are setup and teardown tasks?
Answer: TaskGroups group related tasks in the UI and namespace their IDs; they replaced SubDAGs for visual structure and carry no scheduling overhead. Setup and teardown tasks (@setup, @teardown) model resources that must be created and always cleaned up, such as a temporary cluster or a scratch schema. A teardown runs even if the work between them fails, and in Airflow 3 teardowns also run when a DAG run is terminated early. A teardown failure does not, by default, mark the DAG run failed if the actual work succeeded.
27. How does event-driven scheduling work in Airflow 3?
Answer: You attach an AssetWatcher to an asset. The watcher wraps an event trigger (a trigger class derived from BaseEventTrigger), for example one that reads a message queue through the common messaging provider. When a message arrives, the triggerer creates an asset event, and DAGs scheduled on that asset start. This avoids polling sensors for "start when something happens" workloads.
Production consideration: Airflow still starts a whole DAG run per event, which is fine for "a file landed" or "a batch is ready", not for per-record processing at high volume. For that, keep a streaming system in front and let Airflow react to batch-level events (our Kafka interview questions cover the streaming side).
28. When do you write a custom operator or hook?
Answer: Write a hook when several tasks need to talk to an internal system (a policy admin API, a mainframe gateway) and you want connection handling, auth and retries in one place. Write an operator when a pattern repeats across many DAGs and benefits from templated fields and a consistent interface. Keep TaskFlow functions for one-off logic. In Airflow 3, subclass BaseOperator and BaseHook from airflow.sdk, package them as a normal Python library or internal provider (plugins can no longer register them), and never touch the metadata database directly. If the operator waits on long external work, make it deferrable.
Executors, security and operations
29. Compare the main executors.
Answer:
| Executor | How tasks run | Trade-offs |
|---|---|---|
| LocalExecutor | Subprocesses on the scheduler host | Simple; suits development and small single-machine setups; competes with the scheduler for resources |
| CeleryExecutor | Long-running workers pulling from a broker (Redis or RabbitMQ) | Low start-up latency; you manage workers and the broker; workers share one environment |
| KubernetesExecutor | One pod per task | Strong isolation and per-task resources; pod start-up latency; needs a Kubernetes cluster |
| EdgeExecutor (edge3 provider) | Edge workers in other locations pull work over HTTP | Runs tasks near data or on-premise systems without a direct database or broker link |
| Cloud container executors (for example ECS or Batch, from provider packages) | Each task as a cloud container job | No workers to manage; start-up latency and per-job cost |
Since 2.10 you can configure several executors at once (executor = CeleryExecutor,KubernetesExecutor) and pick one per task or DAG with the executor argument. That replaces the removed hybrid executors: quick tasks on Celery, heavy or dependency-unusual tasks on Kubernetes.
30. KubernetesExecutor versus KubernetesPodOperator: what is the difference?
Answer: The KubernetesExecutor decides where Airflow task code runs: each task instance gets its own pod running the Airflow task runtime with your DAG code. The KubernetesPodOperator is a task that launches an arbitrary container (any language, any image) and watches it; it works with any executor. Use the executor for isolation of Python tasks; use the pod operator to run an existing containerised job, such as a Spark submit or a model-training image, without installing its dependencies into the Airflow image. Both need sensible resource requests, namespaces and service accounts. Our Kubernetes for AI interview questions go deeper on the cluster side.
31. How do secrets backends work, and what is the lookup order?
Answer: Set [secrets] backend to a provider class (for example the AWS Secrets Manager, Azure Key Vault, Google Secret Manager or Vault backends) with backend_kwargs for prefixes and options. When a connection or variable is requested, Airflow searches the configured secrets backend first, then environment variables (AIRFLOW_CONN_*, AIRFLOW_VAR_*), then the metadata database; the order is not configurable, and a workers-specific secrets backend, if set, is checked before the general one. Two gotchas: the UI only shows connections stored in the metadata database, and duplicate keys across backends resolve to the first match, which confuses debugging. Use lookup prefixes and filters so Airflow does not query the secret store for every missing key.
32. How do you monitor an Airflow platform?
Answer: Airflow emits metrics through StatsD or OpenTelemetry. The ones worth dashboards and alerts:
scheduler_heartbeatandscheduler.scheduler_loop_duration: is the scheduler alive and keeping up?dag_processing.total_parse_time,dag_processing.import_errors: parse health.executor.open_slots,executor.queued_tasks,pool.open_slots,pool.starving_tasks: capacity.dagrun.schedule_delay,dagrun.first_task_scheduling_delay,task.queued_duration: latency between "should run" and "is running".
Add remote task logging to object storage or a log platform, health checks for scheduler, triggerer and DAG processor, failure callbacks or notifiers, and deadline alerts for business deadlines. From 3.3, OpenTelemetry timers are recorded as histograms. The general approach is the same as in our AI observability guide: alert on symptoms users feel, investigate with detailed signals.
33. How do you test Airflow DAGs?
Answer: Use layers. First, a DAG integrity test that loads every file and asserts no import errors, plus structural assertions (expected tasks, retries set, owners and tags present). Second, unit tests for task logic, written as plain Python functions you can call without Airflow. Third, run a whole DAG locally with dag.test() or airflow dags test, which executes tasks in one process and is easy to debug. Fourth, integration tests against a staging environment with real connections.
import pytest
from airflow.dag_processing.dagbag import DagBag
@pytest.fixture(scope="session")
def dagbag():
return DagBag(dag_folder="dags/")
def test_no_import_errors(dagbag):
assert dagbag.import_errors == {}
def test_ingest_retries(dagbag):
dag = dagbag.get_dag("claims_ingest")
assert dag.default_args["retries"] == 3
In Airflow 3.3, DagBag is imported from airflow.dag_processing.dagbag and loading a DAG this way needs an initialised metadata database (SQLite is fine in CI).
34. What does a CI/CD pipeline for DAGs look like?
Answer: On every pull request: lint and format; run ruff check --select AIR3 rules to catch removed or deprecated Airflow interfaces; run the DAG integrity and unit tests; and build the image if DAGs need new dependencies. On merge: build and scan a versioned image with pinned Airflow and provider versions (use the official constraints files), deploy DAG code through a Git DAG bundle or image release, and promote from dev to staging to production. Keep connections and variables per environment in the secrets backend, never in the repository. With a Git bundle, every run records the commit it used, which helps audits and rollbacks. Our guide on CI/CD for AI applications covers the GitHub Actions and Argo CD side.
35. Compare self-managed Airflow with managed offerings.
Answer: The main managed options are Amazon MWAA (Managed Workflows for Apache Airflow), Google Cloud's Managed Service for Apache Airflow (previously Cloud Composer), and Astronomer's Astro. All three run the scheduler, API server and database for you and integrate with their cloud's identity, logging and secrets. The trade-offs are less control over versions, configuration and executors, provider choices tied to the platform, a delay before new Airflow releases are offered, and cost models you need to understand. Self-managing, often on Kubernetes with the official Helm chart, gives full control but makes your team responsible for upgrades, database maintenance, scaling and security patches.
Interview tip: Supported Airflow versions on managed platforms change often, and some have restrictions on in-place upgrades from 2.x to 3.x. Say "I would check the provider's supported-versions page" rather than quoting a version.
36. How do you scale Airflow for thousands of DAGs?
Answer: Scale each bottleneck separately. Parsing: run more DAG processor capacity ([dag_processor] parsing_processes, or several DAG processors for different bundles), keep top-level code cheap, and consider longer parse intervals for DAGs that rarely change. Scheduling: run more than one scheduler for high availability and throughput (supported since 2.0, using row-level locks in the database), and tune max_tis_per_query and max_dagruns_to_create_per_loop carefully. Execution: autoscale Celery workers or use Kubernetes. Waiting: deferrable operators and enough triggerers. Database: a properly sized PostgreSQL with connection pooling and regular airflow db clean to purge old rows. Measure first with the metrics in Q32.
If you want hands-on practice building and operating production pipelines like these, Cloudsoft's HORIZON Data Engineering & AI program covers orchestration, cloud data platforms and data for AI, in Ameerpet classrooms or live online.
Airflow for ML and AI pipelines
37. Design an Airflow pipeline that keeps a RAG index fresh.
Answer: Treat the index as an asset with clear ownership and freshness rules.
source docs (SharePoint, S3, wiki)
|
list changed docs (since last run)
|
parse + chunk (mapped per batch)
|
embed (pool limits API rate)
|
upsert to vector store (idempotent ids)
|
delete removed docs from index
|
eval sample queries --> alert if drop
|
outlet: Asset("kb_index")
Key decisions: detect changes incrementally (the 3.3 task state store or a watermark table can hold the last processed cursor); map over batches of documents, not individual chunks; put embedding calls in a pool sized to the model provider's rate limits; use deterministic chunk IDs so retries upsert rather than duplicate; propagate document permissions into metadata; handle deletions; and run a small retrieval evaluation before declaring the index ready. Downstream DAGs, such as a cache warm-up, schedule on the index asset. Our guide to data pipelines for RAG covers ingestion design in depth.
38. How would you orchestrate model retraining with Airflow?
Answer: Trigger retraining on data, not just on time: schedule the training DAG on a curated feature-table asset (optionally combined with a time schedule) so it runs when new labelled data lands. Steps: validate data and check drift; train in an isolated environment (a KubernetesPodOperator with a GPU node pool, or a managed training job through a deferrable provider operator); evaluate against the current production model on a fixed holdout; register the candidate in a model registry such as MLflow; require human approval for promotion if the domain is regulated (a 3.1 approval operator fits here); then deploy through the serving platform's API. Airflow 3's support for runs with no logical date makes ad hoc retrains and hyperparameter sweeps natural. Airflow orchestrates; it should not be the training runtime. The MLOps interview questions guide covers registry, serving and monitoring.
39. Where do human-in-the-loop operators fit in AI workflows?
Answer: Use them where a person must decide before the pipeline continues: approving a model promotion, reviewing low-confidence document extractions, choosing a branch after an LLM-generated summary, or signing off a regulatory report. The task pauses until a user with the right role responds in the UI or through the API, and from 3.2 the UI keeps a history of approvals and rejections. They are suited to decisions measured in minutes or hours, not to chat-speed interaction. For broader patterns see our guide to human-in-the-loop AI.
40. When would you not choose Airflow?
Answer: When the job is really streaming or per-event processing (use a stream processor), when a workflow needs sub-second step latency, when a long-running AI agent needs durable per-step state and complex retries inside a single request (a durable execution engine fits better; see durable, long-running agents), or when a platform already provides adequate built-in scheduling for a simple job, such as the job scheduler in a lakehouse platform. Airflow fits cross-system batch orchestration with dependencies, retries, backfills and audit history. In an interview, naming the cases where it does not fit makes your recommendation more credible.
Real-world scenario questions
41. A DAG's tasks sit in "queued" for a long time and never start. How do you debug it?
Answer: Queued means the scheduler handed the task to the executor but no worker has started it. The problem is almost always between the executor and the workers.
What I would check:
- Executor metrics:
executor.open_slotsandexecutor.queued_tasks. Zero open slots means a capacity problem. - Workers: are Celery workers up and listening on the task's
queue? A task sent to a queue no worker consumes waits forever. On Kubernetes, are pods Pending because of resource requests, quotas, node selectors or image pull errors? - The broker (Celery): is Redis or RabbitMQ healthy and reachable?
- Multiple executors: is the task assigned to an executor that is configured on the scheduler and workers?
- Worker-to-API-server connectivity in Airflow 3: workers must reach the API server's execution API. A network policy or bad URL leaves tasks unable to start.
- Edge executor: are edge workers registered and polling the right queues?
Production consideration: [scheduler] task_queued_timeout (600 seconds by default) makes the scheduler fail or reschedule tasks stuck in queued, so you see a clear error instead of silent waiting. Alert on task.queued_duration rather than waiting for users to notice.
42. Tasks start minutes after their upstream finishes, and the scheduler seems to lag. What do you do?
Answer: Separate parsing lag, scheduling lag and execution capacity, because each has a different fix.
What I would check:
dag_processing.total_parse_timeand per-file parse durations: one slow file with top-level API calls can delay all DAGs. Fix the code first.scheduler.scheduler_loop_durationandscheduler.critical_section_duration: long loops point to database pressure.- Metadata database: CPU, connections, slow queries, and table sizes. Years of task instances and XComs slow everything; run
airflow db cleanwith a retention policy. - Limits:
parallelism, DAGmax_active_tasks, pools. Tasks stuck inscheduledwith free workers mean a limit, not lag. - Scheduler resources: CPU throttling in Kubernetes, or a LocalExecutor competing with the scheduler on the same host.
- Scale out: add a second scheduler, more DAG processor capacity, and spread DAGs that all fire at midnight.
Production consideration: Watch dagrun.first_task_scheduling_delay over time. A slow upward drift usually means the database is growing or DAG count is rising, long before users complain.
43. The platform was down for two days. How do you backfill safely?
Answer: Use Airflow 3's scheduler-managed backfill, after confirming the tasks are idempotent and downstream systems can absorb the load.
What I would check:
- Which DAGs missed runs, and which actually need them. Snapshot-style DAGs may only need the latest run.
- Run a dry run first:
airflow backfill create --dag-id claims_ingest --from-date 2026-09-01 --to-date 2026-09-02 --reprocess-behavior failed --max-active-runs 2 --dry-run. - Choose the reprocess behaviour:
none(only missing runs),failed(missing and failed) orcompleted(rerun everything). - Limit concurrency with the backfill's own
max-active-runsand the pools protecting source systems. - Decide ordering: oldest first if later runs depend on earlier state (and
--run-backwardsis not allowed withdepends_on_past), newest first if users need current data quickly. - Check code version: backfills use the latest bundle version by default; override if the original logic must be reproduced.
- Validate row counts and reconciliation for the backfilled partitions before telling stakeholders.
Production consideration: Backfills can also be started from the UI's Trigger dialog, which is convenient but should still follow the same review. Record who ran it and why.
44. A colleague passes a pandas DataFrame of several hundred MB between tasks through XCom. The UI is slow and the database is growing. What do you change?
Answer: Stop passing data; pass references.
What I would check:
- Find the large XComs (the XCom list in the UI or a query against the XCom table by size) and the DAGs producing them.
- Rewrite tasks to write intermediate data to object storage or a staging table (Parquet on S3, ADLS or GCS) at a deterministic path based on the run, and return only the path.
- Where returning objects is convenient, configure the object storage XCom backend (
XComObjectStorageBackendfrom the common IO provider) with a size threshold, so small values stay in the database and large ones go to object storage with a reference. - Better still, push the transformation into the engine that holds the data (warehouse SQL, Spark) so the worker never loads it.
- Clean up old XCom rows and add a lifecycle policy for the intermediate storage.
Production consideration: Note the 3.3.1 release note that pandas 3 changes how DataFrame XComs are stored and read back. That is one more reason not to rely on serialised DataFrames between tasks.
45. You must migrate a large Airflow 2.x estate to Airflow 3. What is your plan?
Answer: Treat it as a project with an inventory, automated checks, a parallel environment and a phased cut-over.
What I would check:
- Prerequisites: Airflow 2.7 or later (the docs recommend upgrading to the latest 2.x first), supported Python, and a backup of the metadata database.
- Clean up:
airflow db cleanto shrink the database, and fix all DAG parse errors (airflow dags reserializemust run cleanly). - Code: run
ruff check dags/ --select AIR301(breaking changes) and AIR302, then AIR311/AIR312 for recommended updates; apply safe fixes; move imports toairflow.sdkand the standard provider. - Behaviour: find
execution_date, removed context keys,xcom_pullcalls withouttask_ids, SubDAGs, SLAs, and custom operators touching the metadata database. - Scheduling semantics: decide on
create_cron_data_intervalsandcatchupexplicitly before the upgrade. - Config and deployment:
airflow config update --fix, Helm values moved from webserver to API server, a standalone DAG processor, FAB provider if you need FAB auth, plugin views migrated to FastAPI or React. - Run DAGs in a parallel Airflow 3 environment, pause each in the old environment when you switch it, and keep rollback ready.
Real-world example: Consider a GCC team in Hyderabad running a few hundred DAGs for an insurer. Ruff fixes most import changes automatically, but the real work is usually custom operators that queried the metadata database and DAGs that assumed every manual run has a data interval. Migrating those in waves by business domain, with owners signing off each wave, keeps risk contained.
Production consideration: On managed platforms, check whether an in-place upgrade from 2.x to 3.x is supported or whether you must create a new environment and migrate.
46. A worker node was evicted mid-task and the task stayed "running" for a long time. What happened, and how do you handle it?
Answer: The task process died without reporting a final state. Airflow detects this through missed heartbeats: once [scheduler] task_instance_heartbeat_timeout (300 seconds by default) passes, the scheduler treats the task as failed and applies retries.
What I would check:
- Why the node went: spot or preemptible capacity, memory limits (look for out-of-memory kills), or a cluster autoscaler scale-down.
- Whether the external job the task launched is still running. A retry could start a duplicate, so the operator should reattach or the job should be idempotent.
- Pod disruption budgets and safe-to-evict annotations for long tasks on Kubernetes.
- Whether the heartbeat timeout suits your workloads; too short causes false failures under load.
Production consideration: Deferrable operators help here too: if the external job is tracked by a trigger, losing a worker matters less. For Spark, 3.3 added a resumable job mixin to SparkSubmitOperator for surviving worker failures; check the provider docs for the details.
47. Hundreds of sensors waiting for partner files have filled all worker slots, and real work cannot run. Fix it.
Answer: Sensors in poke mode hold a worker slot for the whole wait. Move them off the workers.
What I would check:
- Switch to deferrable versions of the sensors (or
mode="reschedule"where no deferrable version exists) and make sure triggerers are running and sized. - Add timeouts so a missing file fails visibly instead of waiting for days.
- Put sensors in a dedicated pool so they can never take all capacity.
- Consider event-driven scheduling: if the partner's upload can publish a message, an asset watcher can start the DAG when the file arrives, with no sensor at all.
Production consideration: A bank receiving dozens of partner settlement files at different times is a typical case. Asset-triggered DAGs per file type, with a deadline alert for files that are late, tell operations exactly which partner is late.
48. Someone changed a DAG while a long run was in progress, and later tasks in that run used the new logic. How do you prevent this?
Answer: This happens when DAGs come from an unversioned source. With LocalDagBundle or the S3 and GCS bundles, tasks always run the latest code on disk, so a mid-run change affects the remaining tasks.
What I would check:
- Move DAG deployment to a
GitDagBundle, so each run records the commit it started with and every task in that run uses that version. - Decide the rerun policy with
rerun_with_latest_version(3.3): whether clears and backfills should use the original or latest code. - Use the version-aware Graph and Code views to show which version each historical run used during the incident review.
- Require changes to pass CI and deploy via merge, not by editing files on a shared volume.
Production consideration: For audited pipelines (finance, healthcare), "which code produced this report" is a compliance question. A versioned bundle answers it directly.
49. A security review at a bank finds database passwords in DAG files and in task logs. What do you do?
Answer: Rotate first, then fix the design so it cannot happen again.
What I would check:
- Rotate every exposed credential and check access logs for misuse.
- Remove secrets from the repository and its history, and add secret scanning to CI.
- Move connections to a secrets backend (Q31) with least-privilege access for the Airflow identity, and reference only
conn_idin DAGs. - Find why logs contained secrets: values printed by task code, secrets in
bash_commandstrings, or field names not recognised as sensitive. Use the secrets masker, and extend sensitive field names in configuration if needed. - Restrict who can view logs and connections through the auth manager's roles, and in multi-team setups keep connections team-scoped.
- Audit custom operators for direct metadata database access, which Airflow 3 blocks on workers anyway.
Production consideration: Give workers cloud identities (IAM roles, workload identity) instead of static keys wherever the provider supports it, so there is less to leak. Our privacy engineering interview questions cover masking and data protection in more depth.
50. A Bengaluru GCC wants one Airflow platform for twelve data and ML teams. How do you design it?
Answer: Choose between separate deployments per team and a shared deployment with isolation, then standardise everything else.
What I would check:
- Isolation needs: if teams handle regulated data with different access rules, separate deployments (or managed environments) are simpler to reason about. If not, Airflow 3.2's multi-team mode gives per-team DAG bundles, connections, variables, pools and executors in one deployment, but it is experimental, so pilot it first.
- One DAG bundle per team, each from the team's Git repository, with a shared CI template (ruff AIR rules, integrity tests, ownership tags).
- Executors: Celery for quick tasks, Kubernetes for heavy or GPU tasks, per-team queues or namespaces for cost attribution.
- Guardrails: cluster policies to enforce owners, retries, timeouts and pools; a review process for new connections.
- Shared observability: per-team dashboards from executor and pool metrics, and deadline alerts for each team's critical pipelines.
- An upgrade cadence and a provider version policy owned by the platform team.
Production consideration: Asset events cross team boundaries. The docs warn that the permission to create asset events effectively lets a user trigger every DAG scheduled on that asset, so asset access control matters in a shared platform.
Key takeaways
- Airflow orchestrates; heavy data processing belongs in the warehouse, Spark or a container job, and only references should travel through XCom.
- Know the Airflow 3 vocabulary: assets, Task SDK and Task Execution API, API server, DAG processor, DAG bundles, deadline alerts.
- Diagnose by task state: scheduled means limits, queued means executor or workers, deferred means the triggerer.
- Idempotent tasks make retries, reruns and backfills safe; without idempotency, every recovery is risky.
- Deferrable operators, pools and sensible concurrency limits protect worker capacity and fragile downstream systems.
- Migrations to Airflow 3 succeed with automated ruff checks, explicit scheduling decisions and a parallel environment.
- For ML and RAG pipelines, Airflow schedules on data assets and coordinates training, embedding and evaluation steps that run elsewhere.
Interview preparation checklist
- Install Airflow 3 locally with the official constraints file and run the API server, scheduler, DAG processor and triggerer.
- Write a TaskFlow DAG with dynamic task mapping, retries and an asset outlet, and a second DAG scheduled on that asset.
- Convert one sensor to deferrable mode and watch it move to the triggerer.
- Write a DAG integrity test and run a DAG with
dag.test(). - Run a backfill dry run from the CLI and a real backfill from the UI.
- Take an old Airflow 2 DAG and fix it with
ruff check --select AIR301 --fix. - Configure a secrets backend (or environment-variable connections) and confirm nothing sensitive appears in logs.
- Set up StatsD or OpenTelemetry metrics and build a small dashboard of scheduler and pool metrics.
- Prepare two stories from your experience: one incident you debugged and one pipeline you made idempotent.
- Revise SQL and warehouse loading patterns, since most Airflow roles are data engineering roles.
FAQ
Is Airflow still worth learning in 2026?
Yes. Airflow remains widely used for batch data and ML orchestration, and Airflow 3 modernised its architecture, UI and scheduling. Many data engineering job descriptions list it alongside SQL, Python, Spark and a cloud platform.
Should I prepare for Airflow 2 or Airflow 3 questions?
Prepare for Airflow 3, but know the Airflow 2 names. Many companies are still migrating, so interviewers may ask about datasets, SLAs or execution_date and expect you to explain what replaced them.
What skills are required for an Airflow data engineer role?
Strong Python and SQL, DAG design with TaskFlow, retries and idempotency, one cloud platform, a warehouse or lakehouse, containers and Kubernetes basics, Git-based CI/CD, and the ability to debug scheduling and capacity problems.
Do I need to know Kubernetes for Airflow interviews?
For platform and senior roles, usually yes, because many deployments use the Helm chart, the KubernetesExecutor or the KubernetesPodOperator. For DAG-authoring roles, understanding pods, resources and images at a basic level is usually enough.
Can a fresher get an Airflow role?
Freshers are usually hired into data engineering roles where Airflow is one skill among several. A project with a few well-tested DAGs, an asset-triggered pipeline and a written explanation of design choices helps a fresher stand out.
How long does it take to prepare for an Airflow interview?
It depends on your background. An engineer who already writes Python and SQL can cover the Airflow-specific topics and build a project in a few focused weeks. Most of that time should go into building, breaking and debugging DAGs.
Is managed Airflow experience valued as much as self-hosted?
Both are valued. Managed platforms like MWAA, Google Cloud's managed Airflow and Astro are common in enterprises. Self-hosted experience shows deeper operational knowledge, so explain what the platform handled and what you handled yourself.
Are Airflow certifications useful?
Vendor certifications can help a resume pass screening, but interviewers judge hands-on reasoning. Clear answers to scenario questions and a project you can explain in depth usually carry more weight than a certificate alone.
How does Airflow relate to AI and GenAI work?
Airflow orchestrates the batch side of AI systems: document ingestion and embedding for RAG, feature pipelines, model retraining, batch inference and evaluation runs. Airflow 3 added features aimed at these workloads, including runs without a logical date and human-in-the-loop approvals.
Interviewers look for engineers who can keep pipelines correct when things go wrong, not just write DAGs. If you want structured practice with orchestration, cloud data platforms and data for AI, explore the HORIZON program for data engineering and AI. If your goal is broader AI, ML, cloud and security work, the APEX AI, ML, Cloud and Cyber Security program may suit you better. Both run in Ameerpet classrooms or live online; call +91 96660 19191 for a free demo.



