Apache Airflow

0.0(0)
Studied by 0 people
call kaiCall Kai
learnLearn
examPractice Test
spaced repetitionSpaced Repetition
heart puzzleMatch
flashcardsFlashcards
GameKnowt Play
Card Sorting

1/55

encourage image

There's no tags or description

Looks like no tags are added yet.

Last updated 7:29 PM on 9/4/26
Name
Mastery
Learn
Test
Matching
Spaced
Call with Kai
Chat

No analytics yet

Send a link to your students to track their progress

56 Terms

1
New cards
What is Airflow?
A workflow orchestrator. You define DAGs of tasks in Python; the scheduler picks them up, the executor runs them on workers, the metadata DB tracks state, and the UI surfaces it.
2
New cards
What are the benefits of Airflow over cron?
Airflow records every task instance, its state, retries, logs, and lineage — enabling backfills, retries on failure, "rerun from this task," and full observability. Cron has none of that.
3
New cards
What is a DAG?
A Directed Acyclic Graph: a set of tasks plus their dependencies. "Acyclic" means no cycles — a task cannot depend on itself transitively.
4
New cards
What are the three parts of a DAG file?
(1) default_args shared by tasks, (2) the DAG object itself with DAG parameters, (3) the tasks plus their dependencies.
5
New cards
What are DAG default_args?
Per-task defaults applied to every task in the DAG. Each task inherits them unless it overrides at the task level. They control how individual tasks behave: ownership, retries, alerting, resource routing.
6
New cards
What are DAG parameters?
Keyword arguments passed to DAG(...). They control how the DAG as a whole is identified, scheduled, and paced: dag_id, schedule, catchup, concurrency caps, UI metadata.
7
New cards
What is the distinction between default_args and DAG parameters?
default_args define per-task behavior (retries, owner, retry_delay, pool). DAG parameters define DAG-level behavior (schedule, catchup, max_active_runs, dagrun_timeout).
8
New cards
Why must start_date be a static date, never datetime.now()?
A moving start_date shifts every time the DAG file is parsed, so the scheduler cannot compute deterministic logical dates — produces missing or duplicate runs.
9
New cards
What does depends_on_past=True do?
A task instance only runs after the previous run's instance of that task succeeded — serializes runs and blocks the chain on a single failure.
10
New cards
What does the pool default_arg do?
Caps concurrent tasks across all DAGs that share the pool. Used to throttle expensive external resources like an API quota or a DB connection limit.
11
New cards
What does catchup=True do?
When you deploy a DAG with start_date in the past, Airflow runs every missed interval between start_date and now. Default to False for new DAGs (a year-old DAG would kick off 365 runs); set True deliberately for replays/backfills.
12
New cards
What does max_active_runs control?
Cap on concurrent DAG runs for this DAG. Setting to 1 forces strict serial scheduling.
13
New cards
What does dagrun_timeout do?
Kills a DAG run that exceeds the given duration — DAG-level timeout, distinct from per-task retry/timeout settings.
14
New cards
What is an operator?
A class that defines what kind of work a task does — e.g., PythonOperator runs a Python function, BashOperator runs a shell command, SQLExecuteQueryOperator runs SQL.
15
New cards
What is a task?
A configured instance of an operator inside a DAG, with a task_id and concrete parameters. One DAG can have many tasks of the same operator type.
16
New cards
What is a TaskInstance?
A specific run of a task at a specific logical date (e.g., the 'transform' task's run for 2026-04-25). The metadata DB tracks state per TaskInstance; retries produce new TaskInstance attempts, not new Tasks.
17
New cards
Operator vs Task vs TaskInstance?
Operator = the class (what kind of work). Task = a configured node in a DAG. TaskInstance = a specific run of that task at a specific logical date.
18
New cards
Are tasks stateful?
No — tasks are stateless from Airflow's perspective. All metadata lives in the metadata DB. A retry can land on a different worker, so anything written to local disk is gone.
19
New cards
Why does task_id matter?
It is the addressing key — must be unique within the DAG. It is what >> references, what the UI displays, and what logs are filed under. Renaming it is a breaking change: Airflow treats it as a brand-new task with no history.
20
New cards
PythonOperator vs KubernetesPodOperator?
KubernetesPodOperator runs each task in its own pod with its own image and resource limits — no dependency conflicts across tasks, scales horizontally, and a runaway task cannot take down the whole worker.
21
New cards
What is the modern unified SQL operator?
SQLExecuteQueryOperator from the common-sql provider. The conn_id resolves to a connection record whose conn_type (postgres, snowflake, etc.) selects the dialect. PostgresOperator/MySqlOperator are deprecated.
22
New cards
What is EmptyOperator used for?
A placeholder operator used as a fan-in/fan-out anchor or visual marker. Replaced the deprecated DummyOperator in Airflow 2.4+.
23
New cards
When use PythonVirtualenvOperator?
When a Python task needs different dependencies than the Airflow process itself. Caveat: args, kwargs, and the return value are all pickled — non-pickleable objects (DB connections, open file handles) fail at runtime, not at parse.
24
New cards
What is a sensor?
An operator that waits for a condition. It re-checks ("pokes") on an interval until the condition holds or it times out — e.g., FileSensor waits for a file at a path.
25
New cards
What does sensor mode='poke' do?
The default. The sensor holds a worker slot the entire time it waits. Fine for short waits (<5 min); will starve workers for long ones.
26
New cards
What does sensor mode='reschedule' do?
Releases the worker slot between pokes. The task is rescheduled and re-runs on the next interval. Use for any wait of more than a few minutes.
27
New cards
What is a deferrable sensor / operator?
A variant that hands the wait off to the triggerer process (asyncio coroutines). The worker is freed entirely until the condition fires. Best option for waits >10 minutes.
28
New cards
What is ExternalTaskSensor?
A sensor that waits for a specific task in another DAG to succeed for the matching execution date — used to chain DAGs.
29
New cards
What is XCom?
Cross-Communication: the mechanism for passing small data between tasks. Default backend stores values in the metadata DB.
30
New cards
Can you pass large data through XCom?
No. The default backend stores values in the metadata DB; large payloads bloat the DB, slow the UI, and can OOM the scheduler. Pass a pointer (S3/GCS path) instead.
31
New cards
How do you push/pull XCom without the TaskFlow API?
Inside a task receiving **ctx, call ctx['ti'].xcom_push(key='k', value=v) and ctx['ti'].xcom_pull(task_ids='upstream', key='k').
32
New cards
What does the TaskFlow API do for XCom?
Decorating functions with @task lets you call them like normal Python — return values become XCom pushes and arguments become XCom pulls automatically. Dependencies and XCom in one expression.
33
New cards
How do you handle large payloads through XCom safely?
Configure a custom XCom backend (subclass BaseXCom) that offloads the value to object storage while keeping a reference in the metadata DB.
34
New cards
What does >> do?
Defines a downstream edge in the DAG. extract >> transform means transform runs after extract. Mirror operator is <<. Method-call equivalents are set_upstream/set_downstream.
35
New cards
How do you fan-out one task to many in parallel?
extract >> [t1, t2, t3]
36
New cards
How do you fan-in many to one?
[t1, t2, t3] >> finalize
37
New cards
What is the default trigger_rule?
all_success — the task runs only if all upstream tasks succeeded.
38
New cards
What trigger_rule do you need after a BranchPythonOperator?
none_failed_min_one_success — the unchosen branch is *skipped* (not failed), and the default all_success would skip the join task. This rule lets the join run after the live branch only.
39
New cards
What does trigger_rule='all_done' do?
The task runs once all upstream finish, regardless of success or failure. Useful for cleanup tasks.
40
New cards
What does trigger_rule='all_failed' do?
The task runs only if all upstream failed — useful for failure-handling/notification tasks.
41
New cards
What is BranchPythonOperator?
An operator whose Python callable returns the task_id (or list of task_ids) to follow. Other branches are skipped.
42
New cards
What schedule values are valid?
Preset string ('@daily'), cron expression ('0 6 * * MON'), timedelta, None (manual only), a custom Timetable class, or a Dataset/Asset list (data-aware).
43
New cards
When does an @daily DAG actually run?
At the *end* of the interval. A daily DAG with start_date=2026-01-01 makes its first run available at 2026-01-02 00:00, covering 2026-01-01's data.
44
New cards
What do data_interval_start / data_interval_end mean?
The exact window the DAG run covers. For a daily DAG running on 2026-01-02, data_interval_start=2026-01-01 and data_interval_end=2026-01-02.
45
New cards
What is {{ ds }} in Jinja templates?
The DAG run's logical date as YYYY-MM-DD. Equal to data_interval_start.date(). It is the *run's* date, not today.
46
New cards
What is the cron format?
`minute hour day-of-month month day-of-week`. Example: '*/15 * * * *' = every 15 min; '0 0 1 * *' = midnight on the 1st of each month.
47
New cards
When use a Timetable instead of cron?
When the schedule logic is too complex for cron — e.g., "weekdays only excluding US holidays."
48
New cards
What is dataset/asset-aware scheduling?
Schedule a DAG to trigger when an upstream DAG produces a specific Dataset (Asset in 3.x). Replaces brittle ExternalTaskSensor chains with explicit data dependencies.
49
New cards
What is a TaskGroup?
Pure UI grouping of tasks within a DAG. Replaces the deprecated SubDagOperator (which spawned heavy child DAG runs). Every task still belongs to the same DAG run.
50
New cards
Why were SubDAGs deprecated?
They spawned child DAG runs, which is heavy and caused scheduling/performance pitfalls. TaskGroups give the visual grouping without the runtime cost.
51
New cards
What is the triggerer?
A separate Airflow process that runs asyncio coroutines (Triggers) on behalf of deferrable operators. Lets long-running waits free the worker entirely.
52
New cards

Triggers vs TriggerDagRunOperator?

"Triggers" = the asyncio coroutines used by deferrable operators. "TriggerDagRunOperator" = an operator that fires another DAG.

53
New cards
What does TriggerDagRunOperator do?
Fires a run of another DAG by dag_id from inside the current DAG. Useful for chaining DAGs imperatively.
54
New cards
Which principle is violated if you write to /tmp in one task and read it in the next?
Tasks are stateless. The retry, or even the next task, may run on a different worker. Pass data via XCom (small) or shared object storage (large).
55
New cards
Which knob would you turn if a DAG is launching 365 historical runs after deploy?
catchup. Setting catchup=False prevents Airflow from running every missed interval between start_date and now.
56
New cards
A sensor is starving your worker pool during a multi-hour wait. What changes?
Switch from mode='poke' to mode='reschedule', or better, use the deferrable variant (deferrable=True) so the wait moves to the triggerer and the worker is freed entirely.