Airflow courseLesson 5 of 5
Airflow course · Lesson 5 of 5
Airflow XCom, Variables, Connections and Secrets
Move small data between tasks with XCom, keep large data out with custom backends, manage config and credentials with Variables, Connections and secrets backends.
On this page
- Running the examples
- XCom push and pull
- What it is
- How it works
- Pitfalls
- In interviews
- Custom XCom backends
- What it is
- How it works
- Pitfalls
- In interviews
- Variables and Connections
- What it is
- How it works
- Pitfalls
- In interviews
- Params
- What it is
- How it works
- Pitfalls
- In interviews
- Connections and secrets backends
- What it is
- How it works
- Pitfalls
- In interviews
- Task state lifecycle
- What it is
- How it works
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
Tasks in a DAG run in separate processes, often on separate machines, so they cannot share Python variables. Airflow gives you four mechanisms for the information that has to flow anyway: XCom for small results passed between tasks, Variables for configuration, Connections for credentials and endpoints, and params for per-run inputs. This lesson covers each, how secrets backends keep credentials out of the database, and the task states you see while all this happens.
Running the examples
The examples use Airflow 3.3 with a metadata database created by airflow db migrate. The first block also configures Airflow’s local-filesystem secrets backend through environment variables (they must be set before Airflow is imported), so the connection and variable examples later work without a real vault. run_dag is the same helper as in earlier lessons.
import contextlib, io, json, os, tempfile, types, uuid
from datetime import datetime, timedelta
from pathlib import Path
# Secrets files that stand in for a real secrets manager
SECRETS_DIR = tempfile.mkdtemp()
Path(SECRETS_DIR, "connections.json").write_text(json.dumps({
"warehouse": {"conn_type": "postgres", "host": "wh.internal", "port": 5432,
"login": "etl_user", "password": "example-password", "schema": "analytics"},
}))
Path(SECRETS_DIR, "variables.json").write_text(json.dumps({"alert_channel": "#data-alerts"}))
os.environ["AIRFLOW__SECRETS__BACKEND"] = "airflow.secrets.local_filesystem.LocalFilesystemBackend"
os.environ["AIRFLOW__SECRETS__BACKEND_KWARGS"] = json.dumps({
"connections_file_path": f"{SECRETS_DIR}/connections.json",
"variables_file_path": f"{SECRETS_DIR}/variables.json",
})
import pendulum
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()))
XCom push and pull
What it is
XCom (“cross-communication”) is a small key-value store, scoped to a DAG run, that tasks use to pass results to each other. Each XCom is identified by dag_id, run_id, task_id, map_index and a key. A task’s return value is stored automatically under the key return_value.
How it works
- Push: return a value (TaskFlow or any operator with
do_xcom_push=True, the default), or callti.xcom_push(key=..., value=...)for extra keys. - Pull: pass a TaskFlow return value into another task (Airflow pulls for you), call
ti.xcom_pull(task_ids=..., key=...), or use it in a template:{{ ti.xcom_pull(task_ids='extract', key='row_count') }}. xcom_pulldefaults tokey="return_value". With a list oftask_idsit returns a list. For a mapped upstream task you get all map indexes unless you passmap_indexes.- Values are serialised (JSON plus Airflow’s serialisers for common types such as datetimes and dataclasses). Arbitrary pickled objects are not accepted.
- In Airflow 3, workers read and write XComs through the API server (the Task Execution API), not directly in the database, and XComs are cleared when a task instance is cleared and rerun.
from airflow.sdk import dag, task
from airflow.providers.standard.operators.bash import BashOperator
@dag(schedule=None, start_date=datetime(2026, 1, 1))
def xcom_demo():
@task
def extract(ti=None) -> dict:
ti.xcom_push(key="row_count", value=120) # an extra, named XCom
return {"path": "s3://lake/orders/dt=2026-03-01/"} # stored as return_value
shell = BashOperator(task_id="shell", bash_command="echo shell-result") # last stdout line
@task
def report(ti=None) -> dict:
return {
"rows": ti.xcom_pull(task_ids="extract", key="row_count"),
"path": ti.xcom_pull(task_ids="extract")["path"],
"both": ti.xcom_pull(task_ids=["extract", "shell"]),
}
templated = BashOperator(
task_id="templated",
bash_command="echo rows={{ ti.xcom_pull(task_ids='extract', key='row_count') }}",
)
e = extract()
[e, shell] >> report()
e >> templated
state, states, values = run_dag(xcom_demo())
for task_id, value in values.items():
print(f"{task_id:10} {value}")
extract {'path': 's3://lake/orders/dt=2026-03-01/'}
report {'rows': 120, 'path': 's3://lake/orders/dt=2026-03-01/', 'both': [{'path': 's3://lake/orders/dt=2026-03-01/'}, 'shell-result']}
shell shell-result
templated rows=120
Note that the explicit pulls in report do not create dependencies. The [e, shell] >> report() line does; without it report could run before the values exist and pull None.
Pitfalls
- Large values. The default backend stores XComs in the metadata database. Data frames and file contents bloat it, slow the UI and can hit size limits. Pass a path, table name or object key instead, or use a custom backend.
- Pulling from a task that is not upstream: you get
Noneor a stale value, depending on timing. - Relying on XCom across DAG runs. By default a pull only looks in the current run.
- Templated pulls return strings unless the DAG renders native objects.
- Pushing secrets. XCom values are visible in the UI.
In interviews
“How do tasks share data?” Expect to explain XCom’s scope and keys, TaskFlow’s automatic push/pull, and the most important rule: XCom is for metadata, not datasets. Strong candidates say what they pass instead (paths, partitions, table names) and mention custom or object-storage backends for the occasional larger value.
Custom XCom backends
What it is
A custom XCom backend changes where and how XCom values are stored. The common reason is to keep large values out of the metadata database: the backend writes the value to object storage and stores only a reference in the database. It can also add compression, encryption or support for extra types.
How it works
- Subclass
BaseXCom(airflow.sdk.bases.xcomin Airflow 3) and overrideserialize_value(called when a value is pushed; returns what goes into the database) anddeserialize_value(called when it is pulled). Optionally overridepurgeto delete the stored object when the XCom is cleared. - Point Airflow at it with
[core] xcom_backend = my_company.xcom.MyBackendon every component (workers, API server, scheduler), and make the module importable everywhere. - You often do not need to write one: the
apache-airflow-providers-common-iopackage shipsXComObjectStorageBackend, configured with[common.io] xcom_objectstorage_path(any fsspec URL such ass3://...),xcom_objectstorage_threshold(values larger than this many bytes go to storage; small ones stay in the database) and optionalxcom_objectstorage_compression.
A small backend that offloads large values to files, tested by calling its two methods directly (exactly how you would unit test one):
from airflow.sdk.bases.xcom import BaseXCom
XCOM_STORE = Path(tempfile.mkdtemp())
class LocalFileXCom(BaseXCom):
"""Values over THRESHOLD bytes are written to files; the database keeps a reference string."""
PREFIX = "xcom-file://"
THRESHOLD = 1024
@staticmethod
def serialize_value(value, *, key=None, task_id=None, dag_id=None, run_id=None, map_index=None):
payload = json.dumps(value)
if len(payload.encode()) < LocalFileXCom.THRESHOLD:
return BaseXCom.serialize_value(value) # small: store inline
path = XCOM_STORE / f"{dag_id}/{run_id}/{task_id}_{key}_{uuid.uuid4().hex[:8]}.json"
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(payload)
return BaseXCom.serialize_value(LocalFileXCom.PREFIX + str(path)) # large: store a reference
@staticmethod
def deserialize_value(result):
value = BaseXCom.deserialize_value(result)
if isinstance(value, str) and value.startswith(LocalFileXCom.PREFIX):
return json.loads(Path(value[len(LocalFileXCom.PREFIX):]).read_text())
return value
stored = lambda v: types.SimpleNamespace(value=v) # stands in for the stored database row
small = LocalFileXCom.serialize_value({"rows": 3}, key="return_value", task_id="t", dag_id="d", run_id="r")
large = LocalFileXCom.serialize_value(list(range(1000)), key="return_value", task_id="t", dag_id="d", run_id="r")
print("small stored as:", small)
print("large stored as:", large.replace(str(XCOM_STORE), "<store>"))
print("round trip:", LocalFileXCom.deserialize_value(stored(small)),
len(LocalFileXCom.deserialize_value(stored(large))), "items")
small stored as: {'rows': 3}
large stored as: xcom-file://<store>/d/r/t_return_value_eedae318.json
round trip: {'rows': 3} 1000 items
The file names include a random suffix, so yours will differ. The provider’s object-storage backend behaves the same way:
from airflow.providers.common.io.xcom.backend import XComObjectStorageBackend
from airflow.sdk.configuration import conf
conf.set("common.io", "xcom_objectstorage_path", f"file://{XCOM_STORE}/objstore") # s3://bucket/xcom in production
conf.set("common.io", "xcom_objectstorage_threshold", "1024")
ref = XComObjectStorageBackend.serialize_value(list(range(1000)), key="k", task_id="t", dag_id="d", run_id="r")
print(ref.replace(str(XCOM_STORE), "<store>")[:60] + "...")
print(len(XComObjectStorageBackend.deserialize_value(stored(ref))), "items read back")
file://<store>/objstore/d/r/t/b25b162b-3e45-47a9-ab8f-e37964...
1000 items read back
In a deployment you would set these in configuration instead, for example AIRFLOW__CORE__XCOM_BACKEND=airflow.providers.common.io.xcom.backend.XComObjectStorageBackend and AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_PATH=s3://my-bucket/xcom.
Pitfalls
- A backend that is importable on workers but not on the API server: the UI cannot show values, or pulls fail. Install it everywhere.
- Never cleaning up stored objects. Implement
purgeor add a lifecycle rule on the bucket. - Treating a custom backend as permission to pass huge data through Airflow. Workers still download every value they pull; a table belongs in the warehouse, not in XCom.
- Changing the backend on a running system: existing references must still be readable.
In interviews
“What if a task really must hand a few megabytes to the next one?” Explain the object-storage XCom backend (threshold, path, compression) or a custom BaseXCom subclass with serialize_value/deserialize_value, and then say the better design is usually to write the data to storage yourself and pass the path.
Variables and Connections
What it is
Variables are global key-value settings (a bucket name, a feature flag, a list of regions) that any DAG can read. Connections describe how to reach an external system: type, host, port, login, password, schema and a JSON extra. Hooks read connections by conn_id, so credentials never appear in DAG code.
How it works
Both can be defined in the UI, with the CLI (airflow variables set, airflow connections add), through the REST API, as environment variables, or in a secrets backend:
AIRFLOW_VAR_<KEY>defines a variable (AIRFLOW_VAR_TARGET_SCHEMA=analytics).AIRFLOW_CONN_<CONN_ID>defines a connection as a URI (postgres://user:pass@host:5432/db) or as JSON.- Variables and connections from environment variables or secrets backends do not appear in the UI list, and cannot be edited there.
Reading them, in Airflow 3 through airflow.sdk:
os.environ["AIRFLOW_VAR_TARGET_SCHEMA"] = "analytics"
os.environ["AIRFLOW_VAR_LIMITS"] = json.dumps({"max_rows": 1000})
os.environ["AIRFLOW_CONN_LAKE"] = "aws://?region_name=eu-west-1"
from airflow.sdk import BaseHook, Connection, Variable
print(Variable.get("target_schema"))
print(Variable.get("limits", deserialize_json=True))
print(Variable.get("not_defined", default="fallback"))
print(Variable.get("alert_channel")) # from the secrets backend
wh = BaseHook.get_connection("warehouse") # from the secrets backend
print(wh.conn_type, wh.host, wh.port, wh.schema, wh.login)
print(BaseHook.get_connection("lake").extra_dejson) # from AIRFLOW_CONN_LAKE
# Build a connection URI for an environment variable (special characters are encoded for you)
print(Connection(conn_id="tmp", conn_type="postgres", host="db.internal", login="etl",
password="p@ss word", port=5432, schema="sales").get_uri())
analytics
{'max_rows': 1000}
fallback
#data-alerts
postgres wh.internal 5432 analytics etl_user
{'region_name': 'eu-west-1'}
postgres://etl:p%40ss%20word@db.internal:5432/sales
In templates, use {{ var.value.target_schema }}, {{ var.json.limits.max_rows }} and {{ conn.warehouse.host }}. Inside tasks in Airflow 3, Variable.get, Variable.set and connection lookups go through the Task Execution API; tasks no longer get a direct database session.
Pitfalls
Variable.getat the top level of a DAG file. Every parse (and every task start) does a lookup, which may be a database or secrets manager call. Read variables inside tasks or with templates, which are rendered only at run time.- Storing secrets in Variables. Use connections or the secrets backend; Airflow masks values of variables whose names contain words such as
password,secret,tokenorapi_key, but naming is not security. - Using Variables to pass data between tasks or runs. That is XCom (within a run) or storage (across runs).
- Environment variable names: the key part is upper-cased (
AIRFLOW_CONN_MY_DBis the connectionmy_db). - Special characters in URI-form connections must be URL-encoded; the JSON form avoids this.
In interviews
“Where do you keep credentials in Airflow?” Connections, ideally resolved from a secrets backend, referenced by conn_id in hooks; never in DAG code or Variables. Also expect “Why is Variable.get in top-level code bad?” (parse-time lookups) and the difference between Variables (global config), params (per-run input) and XCom (per-run data).
Params
What it is
Params are named, typed inputs for a DAG run, declared on the DAG (or a task) with defaults and validation. When someone triggers a run from the UI, the API or the CLI, Airflow builds a form from them and validates the submitted values. Scheduled runs use the defaults.
How it works
- Declare with
params={"name": Param(default, type=..., ...)}; the extra keywords are JSON Schema (enum,minimum,maximum,format,pattern,items), plus UI hints such astitleanddescription. - Read with
{{ params.name }}in templates,params["name"]in a TaskFlow task (def t(params=None)), orcontext["params"]. - Trigger conf (
{"conf": {...}}in the API,--confin the CLI) overrides params, as long as[core] dag_run_conf_overrides_paramsisTrue(the default). The raw conf is still available asdag_run.conf. - A param with default
Noneand a non-null type makes the value required when triggering. - Task-level
paramsoverride DAG-level ones for that task. - Validation happens when the run is created; an invalid run is rejected before any task starts.
from airflow.sdk import DAG, Param, ParamsDict
@dag(
schedule="@daily",
start_date=datetime(2026, 1, 1),
catchup=False,
params={
"mode": Param("incremental", enum=["incremental", "full"], description="Load strategy"),
"batch_size": Param(500, type="integer", minimum=1, maximum=10_000),
"regions": Param(["eu"], type="array", items={"type": "string"}),
},
)
def params_demo():
@task
def plan(params=None, dag_run=None) -> dict:
return {"mode": params["mode"], "batch": params["batch_size"],
"regions": params["regions"], "raw_conf": dag_run.conf}
@task(params={"batch_size": 50}) # task-level default
def small_batches(params=None) -> int:
return params["batch_size"]
plan()
small_batches()
march_1 = pendulum.datetime(2026, 3, 1, tz="UTC")
print(run_dag(params_demo(), logical_date=march_1)[2])
print(run_dag(params_demo(), logical_date=march_1,
run_conf={"mode": "full", "regions": ["eu", "us"]})[2])
required = ParamsDict({"customer_id": Param(None, type="integer")})
try:
required.validate()
except Exception as exc:
print(type(exc).__name__, "-", str(exc).splitlines()[0])
{'plan': {'mode': 'incremental', 'batch': 500, 'regions': ['eu'], 'raw_conf': {}}, 'small_batches': 50}
{'plan': {'mode': 'full', 'batch': 500, 'regions': ['eu', 'us'], 'raw_conf': {'mode': 'full', 'regions': ['eu', 'us']}}, 'small_batches': 50}
ParamValidationError - Invalid input for param customer_id: None is not of type 'integer'
Pitfalls
- Using
dag_run.confdirectly with no defaults or validation; a typo in a key silently becomesNone. Declare params. - Expecting run conf to change what tasks exist. The DAG structure is fixed at parse time; params only change values at run time.
- Templating params into shell commands without quoting (injection risk).
- Forgetting that scheduled runs always use defaults, so a param must have a sensible default unless the DAG is manual-only.
In interviews
“How do you rerun a DAG for one customer only?” Trigger it with conf that overrides a declared param, validated by its schema, and template the param into the task. Contrast params (per run, user supplied) with Variables (global) and XCom (produced by tasks).
Connections and secrets backends
What it is
A secrets backend lets Airflow read connections and variables from an external secrets manager (AWS Secrets Manager or SSM Parameter Store, Google Secret Manager, Azure Key Vault, HashiCorp Vault) instead of its own database. Credentials stay in the system your security team already manages, with its rotation, audit and access policies.
How it works
- Configure one backend with
[secrets] backend(the class path) and[secrets] backend_kwargs(JSON). The classes come from provider packages. - Lookup order: the configured secrets backend first, then environment variables (
AIRFLOW_CONN_*,AIRFLOW_VAR_*), then the metadata database. The first hit wins. - Each backend maps
conn_idand variable keys to secret names with prefixes, for exampleairflow/connections/warehouse. Setting a prefix tonulldisables lookups of that kind, which avoids paying for a network call on every variable lookup when you only keep connections there. - Optional caching (
[secrets] use_cache, off by default, withcache_ttl_seconds) reduces calls to the secrets manager during parsing. - Connections stored in the metadata database are encrypted with the Fernet key (
[core] fernet_key). Rotate it withairflow rotate-fernet-keyafter adding the new key in front of the old one.
# airflow.cfg (or AIRFLOW__SECRETS__BACKEND / AIRFLOW__SECRETS__BACKEND_KWARGS)
[secrets]
backend = airflow.providers.amazon.aws.secrets.secrets_manager.SecretsManagerBackend
backend_kwargs = {"connections_prefix": "airflow/connections", "variables_prefix": null, "region_name": "eu-west-1"}
# HashiCorp Vault instead:
# backend = airflow.providers.hashicorp.secrets.vault.VaultBackend
# backend_kwargs = {"connections_path": "connections", "variables_path": null, "mount_point": "airflow", "url": "https://vault.internal:8200"}
The local-filesystem backend used on this page (airflow.secrets.local_filesystem.LocalFilesystemBackend, with JSON, YAML or .env files) is useful for tests and local development. In production you would also give the Airflow components a cloud identity (an IAM role, a workload identity) so that the backend itself needs no stored password.
Pitfalls
- A secrets backend lookup for every
Variable.getin top-level code, multiplied by every parse: slow parsing and a large secrets-manager bill. Keep lookups at run time, disable variable lookups you do not need, or enable the cache. - Expecting connections from a secrets backend to appear in the UI. They do not; the UI shows only database connections.
- Losing or changing the Fernet key, which makes stored connection passwords unreadable.
- Different components configured with different backends, so the scheduler and workers see different values.
In interviews
“How would you manage credentials for 50 source systems?” A secrets backend (name the one for your cloud), connections referenced by conn_id, the lookup order (backend, environment, database), prefixes to limit lookups, the Fernet key for anything in the database, and an identity-based login for Airflow itself.
Task state lifecycle
What it is
Every task instance moves through a set of states. Reading them correctly is how you debug “why is my task not running?”.
How it works
The normal path is none → scheduled → queued → running → success:
| State | Meaning |
|---|---|
| none | created, dependencies not yet met |
| scheduled | dependencies met; the scheduler decided it should run |
| queued | handed to the executor, waiting for a worker slot |
| running | executing on a worker |
| success / failed | finished |
| up_for_retry | failed with retries left; it waits retry_delay, then goes back to scheduled |
| up_for_reschedule | a sensor in reschedule mode waiting for its next poke |
| deferred | handed to a trigger in the triggerer |
| upstream_failed | cannot run because an upstream task failed (under its trigger rule) |
| skipped | skipped by a branch, short circuit, trigger rule or AirflowSkipException |
| removed | the task no longer exists in the DAG definition |
| restarting | cleared while it was running; it will be retried |
| awaiting_input | human-in-the-loop operators waiting for a response (Airflow 3.1 and later) |
Each attempt increments try_number. Retries are configured with retries, retry_delay, retry_exponential_backoff and max_retry_delay:
from airflow.sdk import TaskInstanceState
print([s.value for s in TaskInstanceState])
@dag(schedule=None, start_date=datetime(2026, 1, 1))
def lifecycle_demo():
@task(retries=2, retry_delay=timedelta(seconds=1))
def flaky(ti=None) -> str:
if ti.try_number < 2:
raise ConnectionError("transient network error") # first try fails -> up_for_retry
return f"succeeded on try {ti.try_number}"
@task(retries=0)
def broken() -> None:
raise ValueError("bad data") # no retries -> failed
@task
def downstream() -> str:
return "never runs"
flaky()
broken() >> downstream()
print(run_dag(lifecycle_demo())[1:])
['removed', 'scheduled', 'queued', 'running', 'success', 'restarting', 'failed', 'up_for_retry', 'up_for_reschedule', 'upstream_failed', 'skipped', 'deferred', 'awaiting_input']
({'broken': 'failed', 'downstream': 'upstream_failed', 'flaky': 'success'}, {'flaky': 'succeeded on try 2'})
The enum lists every state except “no state yet”, which appears as None. flaky failed once, went to up_for_retry, and succeeded on its second try. broken had no retries, so it failed, and downstream became upstream_failed.
Pitfalls
- Tasks stuck in scheduled: concurrency limits (
max_active_tasks, pools,max_active_tis_per_dag) are full, ordepends_on_pastis waiting on a previous run. - Tasks stuck in queued: the executor has no capacity (workers down, Kubernetes quota, Celery broker problems).
[scheduler] task_queued_timeouteventually fails them. - Running tasks that “disappear”: the worker died and the task’s heartbeat stopped; the scheduler marks it failed (and retries it if retries remain).
- Retrying non-retryable errors (bad data) wastes time; raise
AirflowFailExceptionto fail immediately without retries.
In interviews
Be able to list the states and the normal path, explain up_for_retry, up_for_reschedule and deferred, and give a diagnosis for stuck scheduled versus stuck queued tasks. That distinction (scheduler-side limits versus executor capacity) is what interviewers listen for.
Practice questions
A task returns a 500 MB pandas DataFrame and the next task uses it. What is wrong and how do you fix it?
The return value goes to XCom, by default in the metadata database, which bloats it, slows the scheduler and UI, and may fail on size or serialisation. Write the data to object storage or a warehouse table in the first task and return only the path or table name. If occasional larger values are unavoidable, configure the object-storage XCom backend with a threshold.
Does ti.xcom_pull create a dependency between tasks?
No. Only >>, chain and passing TaskFlow return values create dependencies. A pull from a task that is not upstream may run before the value exists and return None.
What two methods does a custom XCom backend override, and where must it be installed?
serialize_value (on push, returns what is stored in the database, often a reference) and deserialize_value (on pull, turns the stored value back into the object), plus optionally purge. Configure it with [core] xcom_backend and install it on every component that reads or writes XComs: workers, API server and scheduler.
In what order does Airflow look up a connection?
The configured secrets backend, then environment variables (AIRFLOW_CONN_<ID>), then the metadata database. The first match wins.
Variables, params or XCom: where does each of these belong? (a) the S3 bucket name for all DAGs, (b) the customer to reprocess in a manual run, (c) the number of rows the extract task wrote.
(a) A Variable, or better a constant or config file if it rarely changes. (b) A param, overridden by the trigger conf and validated. (c) XCom, returned by the extract task.
A task is stuck in queued for 20 minutes. What do you check?
Executor capacity: are workers running and connected (Celery broker, Kubernetes pods pending on quota or image pulls), and is the task’s queue served by any worker? Queued means the scheduler has done its job; the executor has not started the task. Compare with stuck in scheduled, which points at pools and concurrency limits.
Why is Variable.get() at the top of a DAG file a bad idea, especially with a secrets backend?
Top-level code runs on every parse (every 30 seconds per file by default) and when each task starts, so each call becomes a database or secrets-manager request that slows parsing and costs money. Read variables inside tasks or via templates such as {{ var.value.key }}, which are rendered only at run time.
Key takeaways
- XCom passes small, per-run values between tasks; return values are pushed as
return_value, and pulls do not create dependencies. - Keep data out of XCom. When values must be larger, use the object-storage backend or a
BaseXComsubclass withserialize_valueanddeserialize_value. - Variables hold global configuration, Connections hold credentials and endpoints, params hold validated per-run inputs.
- Secrets backends are checked before environment variables and the database; keep lookups out of top-level code.
- Know the task states: scheduled versus queued tells you whether to look at concurrency limits or executor capacity.
Progress is saved in this browser only. No account needed.