Airflow courseLesson 2 of 5
Airflow course · Lesson 2 of 5
Airflow Operators, Hooks, Providers, Branching and Trigger Rules
Use BashOperator and PythonOperator well, write custom operators and hooks, pick provider packages, run pods, and control flow with branching and trigger rules.
On this page
- Running the examples
- BashOperator
- What it is
- How it works
- Pitfalls
- In interviews
- PythonOperator
- What it is
- How it works
- Pitfalls
- In interviews
- Hooks
- What it is
- How it works
- Pitfalls
- In interviews
- Custom operators
- What it is
- How it works
- Pitfalls
- In interviews
- Provider packages
- What it is
- How it works
- Pitfalls
- In interviews
- KubernetesPodOperator
- What it is
- How it works
- Pitfalls
- In interviews
- Branching with BranchPythonOperator
- What it is
- How it works
- Pitfalls
- In interviews
- Trigger rules
- What it is
- How it works
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
Operators are the verbs of Airflow: each one knows how to do one kind of work, such as run a command, call Python, run SQL or start a container. Hooks are the connectors operators use to reach external systems, and provider packages ship both for a given system. This lesson covers the operators you will use daily, how to write your own, and how to control which tasks run with branching and trigger rules.
Running the examples
The runnable examples use Airflow 3.3 with a metadata database created by airflow db migrate. As in the DAG fundamentals lesson, a small helper registers each in-memory DAG and runs it once with dag.test(), returning the run state, each task’s state and each task’s return value. In your own project you would put the DAG in your DAG folder and run python dags/my_dag.py or airflow dags test <dag_id>.
import contextlib, io, os, json, tempfile
import pendulum
from datetime import datetime
from airflow.dag_processing.bundles.manager import DagBundlesManager
from airflow.dag_processing.dagbag import DagBag, sync_bag_to_db
from airflow.models.xcom import XComModel
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, return values)."""
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)
states, values = {}, {}
with create_session() as session:
for ti in dr.get_task_instances(session=session):
key = ti.task_id if ti.map_index < 0 else f"{ti.task_id}[{ti.map_index}]"
states[key] = str(ti.state)
for x in session.query(XComModel).filter_by(dag_id=dag.dag_id, run_id=dr.run_id, key="return_value"):
key = x.task_id if x.map_index < 0 else f"{x.task_id}[{x.map_index}]"
values[key] = XComModel.deserialize_value(x)
return str(dr.state), dict(sorted(states.items())), dict(sorted(values.items()))
BashOperator
What it is
BashOperator runs a shell command in a Bash subprocess on the worker. It is the quickest way to call a CLI tool (dbt run, aws s3 sync, a script), and its bash_command is templated, so commands can include the run’s dates.
How it works
- In Airflow 3 import it from the standard provider:
from airflow.providers.standard.operators.bash import BashOperator(Airflow 2:airflow.operators.bash). - Templated fields are
bash_command,envandcwd. Abash_commandending in.shor.bashis treated as a template file to load. - Exit codes decide the state: 0 is success, the value of
skip_on_exit_code(99 by default) marks the task skipped, any other non-zero code fails it. - The last line of stdout is pushed to XCom as the return value (turn off with
do_xcom_push=False). envreplaces the worker’s environment unless you also setappend_env=True.@task.bashis the TaskFlow version: the function returns the command string, which lets you build it in Python.
from airflow.sdk import DAG, dag, task
from airflow.providers.standard.operators.bash import BashOperator
with DAG("bash_demo", schedule=None, start_date=datetime(2026, 1, 1)) as bash_dag:
BashOperator(
task_id="templated",
bash_command="echo processing dt={{ ds }} for $REGION",
env={"REGION": "eu"},
append_env=True, # keep PATH and the rest of the worker's environment
)
BashOperator(task_id="nothing_to_do", bash_command="echo no new files; exit 99") # skipped
BashOperator(task_id="broken", bash_command="echo about to fail; exit 3") # failed
@task.bash
def build_command(table: str = "orders") -> str:
return f"echo vacuuming {table}"
build_command()
state, states, values = run_dag(bash_dag, logical_date=pendulum.datetime(2026, 3, 1, tz="UTC"))
print(state)
print(states)
print(values)
failed
{'broken': 'failed', 'build_command': 'success', 'nothing_to_do': 'skipped', 'templated': 'success'}
{'build_command': 'vacuuming orders', 'templated': 'processing dt=2026-03-01 for eu'}
The DAG run failed because broken failed, the exit-99 task was skipped rather than failed, and each successful task’s last stdout line became its return value.
Pitfalls
- Shell injection. Never template untrusted input (for example run conf typed by a user) straight into
bash_command. Pass it throughenvand quote it. envwithoutappend_env=TruewipesPATH, so commands that worked locally suddenly are “not found”.- Long-running heavy jobs on the worker. Bash runs on the Airflow worker; for heavy work submit to an external system or a pod.
- A literal command that ends in
.shis treated as a template path. Add a trailing space or settemplate_exthandling deliberately. - Relying on the working directory. Each run gets a temporary directory unless you set
cwd.
In interviews
Expect “How does BashOperator decide success?” (exit code, with 99 meaning skip by default) and “How do you pass the run date into a script?” (template {{ ds }} into the command or env). Mentioning the shell-injection risk of templating run conf, and that output goes to XCom from the last stdout line, shows hands-on experience.
PythonOperator
What it is
PythonOperator calls a Python function on the worker. In new code the @task decorator is the usual way to do the same thing, but you will meet PythonOperator in most existing DAGs and interviews still ask about it.
How it works
- Import:
from airflow.providers.standard.operators.python import PythonOperator(Airflow 2:airflow.operators.python). python_callableis the function;op_argsandop_kwargsare its arguments; both are templated, as istemplates_dict.- Context values arrive as keyword arguments if the function declares them (
ds,ti,params,logical_date) or accepts**context. Airflow 1’sprovide_context=Trueis no longer needed or accepted. - The return value is pushed to XCom.
- Related operators in the same module:
BranchPythonOperator,ShortCircuitOperator(skips everything downstream when the callable returns a falsy value),PythonVirtualenvOperatorandExternalPythonOperator(run in a different Python environment, useful for conflicting dependencies).
from airflow.providers.standard.operators.python import PythonOperator, ShortCircuitOperator
def summarise(table, min_rows, ds=None, **context):
run_id = context["run_id"]
return f"{table}: checked {ds} (min_rows={min_rows!r}, run={run_id.split('__')[0]})"
def has_new_data(ds=None):
return ds.endswith("-01") # pretend data only lands on the first of the month
with DAG("python_demo", schedule=None, start_date=datetime(2026, 1, 1)) as python_dag:
gate = ShortCircuitOperator(task_id="has_new_data", python_callable=has_new_data)
check = PythonOperator(
task_id="summarise",
python_callable=summarise,
op_args=["orders"],
op_kwargs={"min_rows": "{{ params.min_rows }}"}, # templated
params={"min_rows": 10},
)
gate >> check
print(run_dag(python_dag, logical_date=pendulum.datetime(2026, 3, 1, tz="UTC"))[1:])
print(run_dag(python_dag, logical_date=pendulum.datetime(2026, 3, 2, tz="UTC"))[1:])
({'has_new_data': 'success', 'summarise': 'success'}, {'has_new_data': True, 'summarise': "orders: checked 2026-03-01 (min_rows='10', run=manual)"})
({'has_new_data': 'success', 'summarise': 'skipped'}, {})
On 2 March the short-circuit returned False, so summarise was skipped instead of failing. Notice also that the templated min_rows arrived as the string '10', not the integer 10: Jinja renders to text unless the DAG sets render_template_as_native_obj=True.
Pitfalls
- Heavy imports at the top of the DAG file for use inside the callable. Import inside the function so parsing stays fast.
- Returning large objects (they go to XCom).
- Expecting templated
op_kwargsto keep their types (see above). - Running conflicting library versions in the worker environment; use
PythonVirtualenvOperator,ExternalPythonOperator, or a container.
In interviews
A common question is “PythonOperator or @task?”. They do the same thing; @task is shorter and wires XCom and dependencies automatically, while PythonOperator is explicit and common in older code. Mention ShortCircuitOperator for “skip the rest if there is nothing to do”, and virtualenv-based operators for dependency isolation.
Hooks
What it is
A hook is a reusable client for an external system that knows how to read credentials from an Airflow connection. PostgresHook, S3Hook and HttpHook are examples. Operators use hooks internally; you also call hooks directly inside @task functions when no operator fits.
How it works
A connection is a named record (conn_id) with a type, host, login, password, schema, port and a JSON extra. It can come from the metadata database, an environment variable AIRFLOW_CONN_<CONN_ID> (URI or JSON), or a secrets backend. A hook’s job is to turn that connection into a working client and offer convenience methods (get_records, load_file, run).
You write a custom hook by subclassing BaseHook:
from airflow.sdk import BaseHook
# A connection supplied as an environment variable (JSON form), as you might in CI or a container
os.environ["AIRFLOW_CONN_WEATHER_API"] = json.dumps({
"conn_type": "http", "host": "api.example.com", "schema": "https",
"login": "svc_airflow", "password": "not-a-real-secret", "extra": {"timeout": 5},
})
class WeatherHook(BaseHook):
"""Builds request settings for the weather API from an Airflow connection."""
conn_name_attr = "weather_conn_id"
default_conn_name = "weather_api"
conn_type = "http"
hook_name = "Weather API"
def __init__(self, weather_conn_id: str = default_conn_name):
super().__init__()
self.weather_conn_id = weather_conn_id
def get_conn(self) -> dict:
conn = self.get_connection(self.weather_conn_id)
return {
"base_url": f"{conn.schema}://{conn.host}",
"auth": (conn.login, "***"), # never print real secrets
"timeout": conn.extra_dejson.get("timeout", 30),
}
print(WeatherHook().get_conn())
{'base_url': 'https://api.example.com', 'auth': ('svc_airflow', '***'), 'timeout': 5}
The hook never hard-codes a host or password, so the same DAG works in dev, staging and production with different connections.
Pitfalls
- Creating a hook (and opening a connection) at the top level of the DAG file. Hooks belong inside
executeor a task function. - Hard-coding credentials in DAG code or Variables instead of connections.
- Printing connection objects or URIs in logs. Airflow masks known secret fields, but a password embedded in a custom string may slip through.
- Rebuilding the same client for every row; create it once per task.
In interviews
“Operator vs hook?” is a staple. A hook connects to a system and exposes methods; an operator is a task that does one unit of work, usually by calling a hook. Strong answers mention connections (conn_id, environment variable form, secrets backends) and that hooks keep credentials out of DAG code.
Custom operators
What it is
A custom operator is a subclass of BaseOperator that packages a repeated piece of work (for example “copy this table from Postgres to S3 for this date”) so DAG authors use it like any built-in operator.
How it works
- Subclass
BaseOperator(fromairflow.sdk; Airflow 2:airflow.models.BaseOperator) and implementexecute(self, context). Its return value goes to XCom. - Declare
template_fieldsso arguments such as dates are rendered by Jinja beforeexecuteruns. Templates are not rendered in__init__. - Keep
__init__cheap: it runs at parse time. Create hooks and connections insideexecute. - Raise an exception to fail (and retry); raise
AirflowSkipExceptionto skip. - Optional:
ui_color,template_ext,on_kill()for cleanup when a task is stopped, and deferral for long waits (covered in the sensors lesson).
from airflow.sdk import BaseOperator
from airflow.sdk.exceptions import AirflowSkipException
class ExportTableOperator(BaseOperator):
"""Pretend export of one table partition; shows templating, hooks and skipping."""
template_fields = ("table", "partition", "target_uri")
ui_color = "#e8f4fd"
def __init__(self, *, table: str, partition: str, target_uri: str,
weather_conn_id: str = "weather_api", **kwargs):
super().__init__(**kwargs) # task_id, retries and other BaseOperator args
self.table = table
self.partition = partition
self.target_uri = target_uri
self.weather_conn_id = weather_conn_id
def execute(self, context):
if self.table.startswith("tmp_"):
raise AirflowSkipException(f"{self.table} is a scratch table")
settings = WeatherHook(self.weather_conn_id).get_conn() # hook created at run time
self.log.info("Exporting %s for %s", self.table, self.partition)
return {"table": self.table, "target": self.target_uri, "api": settings["base_url"]}
with DAG("custom_operator_demo", schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False) as op_dag:
ExportTableOperator(task_id="export_orders", table="orders", partition="{{ ds }}",
target_uri="s3://lake/orders/dt={{ ds }}/")
ExportTableOperator(task_id="export_scratch", table="tmp_orders", partition="{{ ds }}",
target_uri="s3://lake/tmp/")
print(run_dag(op_dag, logical_date=pendulum.datetime(2026, 3, 1, tz="UTC"))[1:])
({'export_orders': 'success', 'export_scratch': 'skipped'}, {'export_orders': {'table': 'orders', 'target': 's3://lake/orders/dt=2026-03-01/', 'api': 'https://api.example.com'}})
Package custom operators and hooks as a normal Python package (or under the DAG bundle’s plugins/shared folder) so they can be unit tested and versioned separately from DAG files.
Pitfalls
- Doing I/O in
__init__. It runs every time the file is parsed. - Expecting a templated field to be rendered in
__init__; it is rendered just beforeexecute. - Forgetting
**kwargsandsuper().__init__(**kwargs), which breakstask_id,retriesand the rest. - Writing a custom operator for something a provider already ships. Search the provider first.
In interviews
Interviewers may ask you to sketch one. Show the class with template_fields, a cheap __init__, an execute(context) that builds a hook, logs and returns a small value, and explain idempotency (overwrite the target partition). Bonus points for mentioning on_kill and deferrable operators for long waits.
Provider packages
What it is
Providers are separately installed and versioned packages (apache-airflow-providers-<name>) that contain operators, hooks, sensors, triggers, connection types and sometimes executors for one system: Amazon, Google, Databricks, Snowflake, Postgres, Celery, Kubernetes and many more. Airflow core stays small, and each integration can release fixes without waiting for an Airflow release.
How it works
- Install what you need, pinned with the Airflow constraints file:
pip install "apache-airflow-providers-amazon" --constraint <constraints-url>. - Imports follow the pattern
airflow.providers.<name>.<operators|hooks|sensors|triggers>.<module>. - In Airflow 3 the everyday operators moved out of core into
apache-airflow-providers-standard:BashOperator,PythonOperator,BranchPythonOperator,EmptyOperator,TriggerDagRunOperator,ExternalTaskSensor,FileSensor, the time sensors and more. Old imports such asairflow.operators.bashare gone; the Airflow 3 upgrade guide lists the new paths, and theruffAIR rules can rewrite them. - Executors that need extra infrastructure live in providers too:
CeleryExecutorin the Celery provider andKubernetesExecutorin thecncf.kubernetesprovider.
You can list what is installed:
from airflow.providers_manager import ProvidersManager
for name, info in sorted(ProvidersManager().providers.items()):
print(f"{name:42} {info.version}")
apache-airflow-providers-common-compat 1.20.0
apache-airflow-providers-common-io 1.10.0
apache-airflow-providers-common-sql 2.2.0
apache-airflow-providers-smtp 3.1.0
apache-airflow-providers-standard 1.20.0
From the command line, airflow providers list shows the same.
Pitfalls
- Upgrading Airflow without upgrading providers (or the reverse). Check each provider’s minimum Airflow version and install with the constraints file.
- Copying Airflow 2 import paths into Airflow 3 code. Many now fail with
ModuleNotFoundError. - Installing every provider “just in case”. Each brings dependencies that can conflict.
In interviews
Expect “What is a provider?” and, for Airflow 3 roles, “What changed with operator imports?”. Explain that integrations are versioned independently of core, that Airflow 3 moved the basic operators into the standard provider, and that you pin providers with the constraints file.
KubernetesPodOperator
What it is
KubernetesPodOperator (KPO), from the apache-airflow-providers-cncf-kubernetes package, starts a pod in a Kubernetes cluster to run any container image, waits for it, streams its logs and marks the task by the container’s exit code. It is the standard way to run work with its own dependencies, language or resources without installing anything on Airflow workers.
How it works
from kubernetes.client import models as k8s
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator
transform = KubernetesPodOperator(
task_id="transform_orders",
name="transform-orders",
namespace="data-jobs",
image="registry.example.com/etl/orders:1.4.2", # pin a tag, never :latest
cmds=["python", "-m", "etl.orders"],
arguments=["--date", "{{ ds }}"], # templated
env_vars={"TARGET": "s3://lake/orders/"},
container_resources=k8s.V1ResourceRequirements(
requests={"cpu": "1", "memory": "2Gi"},
limits={"memory": "4Gi"},
),
service_account_name="etl-runner",
get_logs=True,
on_finish_action="delete_pod", # or keep_pod / delete_succeeded_pod
do_xcom_push=False, # True: reads /airflow/xcom/return.json via a sidecar
kubernetes_conn_id="kubernetes_default", # or in_cluster=True when Airflow runs in the cluster
deferrable=True, # wait in the triggerer instead of holding a worker slot
)
The TaskFlow form is @task.kubernetes(image=..., namespace=...), which runs the decorated function inside the pod.
KPO vs KubernetesExecutor is a classic confusion. The KubernetesExecutor decides where every task runs (each task instance in its own pod running Airflow’s worker code). KPO is a single task that launches your container, and it works with any executor (Local, Celery or Kubernetes).
Pitfalls
- Using
:latestimages, which makes reruns non-reproducible. - No resource requests or limits, so pods get evicted or starve neighbours.
- Leaving finished pods behind (
on_finish_action="keep_pod") and filling the namespace. - Expecting XCom to work without writing
/airflow/xcom/return.jsonin the container and settingdo_xcom_push=True. - Older parameter names:
is_delete_operator_podandresourceswere replaced byon_finish_actionandcontainer_resourcesin recent provider versions; check your provider version.
In interviews
“How would you run a job that needs a different Python version or a Java tool?” KPO with a pinned image, resource requests and limits, a service account, logs streamed to Airflow, and deferrable mode for long jobs. Be ready to contrast it with the KubernetesExecutor.
Branching with BranchPythonOperator
What it is
Branching chooses at run time which downstream path to follow. The branch task returns the task_id (or list of ids) to run; every other directly downstream task is skipped, and skips propagate down those paths.
How it works
- TaskFlow:
@task.branch. Classic:BranchPythonOperator(python_callable=...). Others in the family:BranchDateTimeOperator,BranchDayOfWeekOperator, andBaseBranchOperatorfor your own. - Return a task id, a list of ids, or
Noneto skip all downstream tasks. For a task inside a group, return the full id (group.task). - The join task after the branches needs a trigger rule that tolerates skipped parents, usually
none_failed_min_one_success. With the defaultall_successthe join is skipped too, because one parent was skipped.
from airflow.sdk import TriggerRule
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.providers.standard.operators.python import BranchPythonOperator
@dag(schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False)
def branching_demo():
@task.branch
def choose_load(ds=None) -> str:
return "full_reload" if ds.endswith("-01") else "incremental_load" # full reload monthly
full = EmptyOperator(task_id="full_reload")
incremental = EmptyOperator(task_id="incremental_load")
join_default = EmptyOperator(task_id="join_all_success") # default rule
join_right = EmptyOperator(task_id="join_none_failed",
trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS)
choose_load() >> [full, incremental]
[full, incremental] >> join_default
[full, incremental] >> join_right
branch_dag = branching_demo()
print(run_dag(branch_dag, logical_date=pendulum.datetime(2026, 3, 1, tz="UTC"))[1])
print(run_dag(branch_dag, logical_date=pendulum.datetime(2026, 3, 2, tz="UTC"))[1])
{'choose_load': 'success', 'full_reload': 'success', 'incremental_load': 'skipped', 'join_all_success': 'skipped', 'join_none_failed': 'success'}
{'choose_load': 'success', 'full_reload': 'skipped', 'incremental_load': 'success', 'join_all_success': 'skipped', 'join_none_failed': 'success'}
The classic form is identical in behaviour:
with DAG("branch_classic", schedule=None, start_date=datetime(2026, 1, 1)) as classic_branch:
pick = BranchPythonOperator(task_id="pick", python_callable=lambda: ["eu", "us"]) # several paths
paths = [EmptyOperator(task_id=r) for r in ("eu", "us", "apac")]
pick >> paths
print(run_dag(classic_branch)[1])
{'apac': 'skipped', 'eu': 'success', 'pick': 'success', 'us': 'success'}
Pitfalls
- The join task skipped under
all_success, the most common branching bug. Usenone_failed_min_one_success. - Returning a task id that is not directly downstream of the branch task; Airflow raises an error.
- Branch logic based on the wall clock instead of the run’s dates, so backfills take the wrong path.
- Using a branch when you only want “stop if there is nothing to do”;
ShortCircuitOperatoror@task.short_circuitis simpler.
In interviews
Be ready to draw a branch and a join and name the join’s trigger rule. Explain that unchosen paths are skipped (not failed) and that skips propagate under all_success. Mention @task.branch as the modern form.
Trigger rules
What it is
A trigger rule decides when a task may run based on the states of its direct upstream tasks. The default, all_success, runs a task only if every parent succeeded. Other rules let you build cleanup tasks, alerts and joins.
How it works
| Rule | Runs when | Typical use |
|---|---|---|
all_success (default) |
every parent succeeded | normal pipelines |
all_failed |
every parent failed or was upstream_failed | fallback path |
all_done |
every parent finished, whatever the state | cleanup, teardown |
all_done_min_one_success |
all parents done and at least one succeeded | partial-success reporting |
all_skipped |
every parent was skipped | “nothing ran” notice |
one_success |
as soon as one parent succeeded | first-of-several sources |
one_failed |
as soon as one parent failed | early alert |
one_done |
as soon as one parent succeeded or failed | race patterns |
none_failed |
no parent failed (skips are fine) | joins after a branch |
none_failed_min_one_success |
no parent failed and at least one succeeded | the usual branch join |
none_skipped |
no parent was skipped | require every path to have run |
always |
immediately, regardless of parents | rarely the right choice |
States that matter: upstream_failed is set on a task whose all_success condition can never be met because a parent failed; skipped propagates through all_success too. Setup and teardown tasks (.as_setup(), .as_teardown()) are a more structured alternative to all_done cleanup: a teardown runs even when the work fails, and does not hide the failure from the DAG run state.
@dag(schedule=None, start_date=datetime(2026, 1, 1))
def trigger_rules_demo():
@task
def extract():
return "ok"
@task
def transform():
raise ValueError("bad input file")
@task
def load():
return "loaded"
@task(trigger_rule=TriggerRule.ONE_FAILED)
def alert():
return "paged on-call"
@task(trigger_rule=TriggerRule.ALL_DONE)
def cleanup():
return "temp files removed"
@task(trigger_rule=TriggerRule.NONE_FAILED)
def publish():
return "published"
e, t = extract(), transform()
l = load()
[e, t] >> l
[e, t] >> alert()
[e, t] >> cleanup()
l >> publish()
state, states, values = run_dag(trigger_rules_demo())
print(state)
for task_id, s in states.items():
print(f"{task_id:10} {s}")
failed
alert success
cleanup success
extract success
load upstream_failed
publish upstream_failed
transform failed
transform failed, so load became upstream_failed and that propagated to publish. alert and cleanup still ran because their rules allow a failed parent. The DAG run is failed because a task that is a leaf (publish) did not succeed.
Pitfalls
all_donecleanup tasks as leaves can make a DAG run look successful even though work failed, because a DAG run’s state comes from its leaf tasks. Prefer teardown tasks, or keep a normal task as a leaf.alwaysignores dependencies, which is almost never what you want.- Trigger rules look only at direct parents, not the whole upstream graph.
- Using
one_successto “speed up” a join, then reading output that the slower parent has not produced yet.
In interviews
Expect “How do you make a cleanup task run even if something failed?” (all_done, or a teardown task) and “Which rule for joining after a branch?” (none_failed_min_one_success). Explaining upstream_failed and how a DAG run’s final state is derived from leaf tasks shows depth.
Practice questions
What is the difference between an operator, a task and a hook?
An operator is a class describing a kind of work (run Bash, run SQL). A task is an instance of an operator (or a @task function) placed in a DAG with a task_id. A hook is a client for an external system that reads an Airflow connection; operators use hooks to do their work, and you can call hooks directly inside tasks.
Your BashOperator task should be marked skipped, not failed, when there is no input file. How?
Exit with code 99 (the default skip_on_exit_code), or set skip_on_exit_code to another code your script uses. Any other non-zero exit fails the task.
After a branch, the join task is always skipped. Why, and how do you fix it?
The join uses the default all_success rule, and one of its parents (the branch not taken) is skipped, so the skip propagates. Set the join’s trigger_rule to none_failed_min_one_success (or none_failed).
Where should a custom operator create its database connection, and why?
Inside execute(), through a hook. __init__ runs every time the DAG file is parsed, which happens constantly on the DAG processor, so connecting there would open connections at parse time, slow parsing and possibly fail imports. Templated fields are also only rendered just before execute.
You upgraded to Airflow 3 and from airflow.operators.bash import BashOperator fails. Why?
In Airflow 3 the standard operators were moved into the apache-airflow-providers-standard package. Import from airflow.providers.standard.operators.bash and make sure that provider is installed.
KubernetesPodOperator or KubernetesExecutor?
They solve different problems. The KubernetesExecutor runs every Airflow task instance in its own pod, which is a deployment choice for the whole installation. KubernetesPodOperator is a single task that runs a container image you choose, with any executor. Use KPO when one task needs its own image, dependencies or resources.
A cleanup task with trigger_rule=“all_done” makes failed runs show as successful. Why?
A DAG run’s final state is based on its leaf tasks. If the only leaf is a cleanup task that succeeds, the run can be marked successful even though upstream work failed. Mark the cleanup as a teardown task (.as_teardown()), which Airflow ignores when deciding the run state, or make sure a normal task is also a leaf.
Key takeaways
- BashOperator decides state by exit code (99 means skip by default) and pushes its last stdout line to XCom; keep untrusted input out of the command.
- PythonOperator and
@taskdo the same job; templatedop_kwargsrender to strings unless native rendering is on. - Hooks turn connections into clients; operators use hooks; custom operators keep
__init__cheap and work inexecute. - Providers version integrations separately from core, and in Airflow 3 the basic operators live in
apache-airflow-providers-standard. - KubernetesPodOperator runs your container as one task and is not the same thing as the KubernetesExecutor.
- Branch with
@task.branchand join withnone_failed_min_one_success; use trigger rules and teardown tasks deliberately, because they change how failures show up.
Progress is saved in this browser only. No account needed.