Part 8
Sensors & Waiting
46HttpSensor#
What it does: Polls an HTTP endpoint until it returns a 2xx status, with an optional response-body regex check.
Parameters:
| Param | Type | Default | Notes |
|---|---|---|---|
endpoint | string (URL) | required | |
connection_id | string | none | If set, resolves auth (bearer token/basic auth) from a stored HTTP connection. |
headers | dict | {} | Merged on top of anything resolved from connection_id, or used on its own. |
response_check_regex | regex string | none | Uses Go's RE2 regex syntax (no lookahead/lookbehind/backreferences). |
poke_interval | integer (seconds) | 60 | |
timeout | integer (seconds) | 3600 | Cumulative across the whole wait (see feature 49). |
mode | string | "poke" | See feature 50 — HttpSensor's default differs from the other three sensors. |
soft_fail | boolean | False | Timeout becomes skipped instead of failed. |
How it works:
- Set
connection_idto reuse stored credentials for the target endpoint; addheaders=for anything not covered by the connection, or to override it. - A network/DNS error is treated the same as "condition not met" — it keeps retrying at
poke_intervaluntiltimeout, rather than failing fast. response_check_regexonly gets evaluated once a 2xx response is actually received — a regex syntax error surfaces at that point, not before.- Response body is capped at 1MB for the regex check.
- Defaults to
mode="poke", unlike the other three sensors — setmode="reschedule"explicitly if you want it to free its worker slot between checks.
Example:
from dag_parser.dynamic.operators import HttpSensor
wait_for_api = HttpSensor(
task_id="wait_for_api",
endpoint="https://api.example.com/v1/status",
connection_id="internal_api", # resolves auth from the stored connection
response_check_regex=r'"status"\s*:\s*"ready"',
poke_interval=30,
timeout=1800,
mode="reschedule", # override the poke default explicitly
soft_fail=True, # timeout -> skipped, not failed
)47SqlSensor#
What it does: Runs a SQL query on every poke and succeeds once the result has at least one row whose first column is truthy.
Parameters:
| Param | Type | Default | Notes |
|---|---|---|---|
sql | string | required | |
connection_id | string | required | The external database connection to run the query against. |
poke_interval | integer (seconds) | 60 | |
timeout | integer (seconds) | 3600 | Cumulative (feature 49). |
mode | string | "reschedule" |
How it works:
connection_iddetermines which external database the query runs against — point it at any registered connection whose type supports a SQL query (Snowflake, Postgres, MySQL, Redshift, etc.).- If you need to wait on another DAG/task's state specifically,
ExternalTaskSensor(feature 48) is the purpose-built tool for that instead of querying tables directly. - Truthiness of the first cell:
NULL,0,"","0", and case-insensitive"false"are falsy; anything else (including a query returning zero rows, which behaves the same as a falsy first column) is not a match. - A malformed SQL string is a genuine error, not "condition not met" — it's retried each interval until timeout, then surfaces as a failure.
Example:
from dag_parser.dynamic.operators import SqlSensor
wait_for_flag = SqlSensor(
task_id="wait_for_flag",
sql="SELECT ready_flag FROM control.batch_status WHERE batch_date = '{{ ds }}'",
connection_id="warehouse_snowflake",
poke_interval=60,
timeout=3600,
mode="reschedule",
)48TimeSensor#
What it does: Waits until the current wall-clock time reaches or passes a target time of day.
Parameters:
| Param | Type | Default | Notes |
|---|---|---|---|
target_time | string, "HH:MM" or "HH:MM:SS" | required | Evaluated in the DAG's own timezone (feature 3). |
poke_interval | integer (seconds) | 60 | |
timeout | integer (seconds) | 3600 | Cumulative (feature 49). |
mode | string | "reschedule" |
How it works:
target_timeis evaluated against the DAG's owntimezonesetting (feature 3) — a DAG declaredtimezone="America/New_York"withtarget_time="09:00"waits for 09:00 US Eastern, whatever timezone the underlying infrastructure runs in.- Always targets today's date — a backfill run representing a historical date still waits for
target_timeon whatever day it actually executes. - If the condition is already true when the sensor first checks (e.g. it's 14:00 and
target_time="09:00"), it succeeds immediately — there's no "wait until tomorrow" behavior. - Defaults to
mode="reschedule", which is efficient for long waits (e.g. "wait until market open") since it only claims a worker slot for each brief check.
Example:
from dag_parser.dynamic.operators import TimeSensor
with DAG(
dag_id="regional_market_open",
schedule="@daily",
timezone="America/New_York",
start_date=datetime(2026, 1, 1),
) as dag:
wait_until_market_open = TimeSensor(
task_id="wait_until_market_open",
target_time="09:30:00", # 09:30 US Eastern, matching the DAG's timezone
poke_interval=60,
timeout=3600 * 2,
mode="reschedule",
)49ExternalTaskSensor#
What it does: Waits for a task (or an entire DAG run) in another DAG to reach one of a set of allowed states.
Parameters:
| Param | Type | Default | Notes |
|---|---|---|---|
external_dag_id | string | required | |
external_task_id | string | none | Omit to wait on the whole DAG run's own state instead of one task. |
allowed_states | list of strings | ["success"] | Case-insensitive match. |
failed_states | list of strings | ["failed"] | Case-insensitive match; checked first, ends the wait as a failure. |
execution_date | string (ISO or templated) | latest run | Pin to a specific external run rather than always polling the newest one. |
poke_interval / timeout / mode | — | 60 / 3600 / "reschedule" | Same semantics as the other sensors. |
How it works:
- Without
execution_date, it always resolves to the latest run of the external DAG/task — if your DAG runs more often than the external one, every check ends up polling the same external run. Passexecution_dateexplicitly if you need a specific, matched run. - If the external state matches neither
allowed_statesnorfailed_states(e.g. it's stillrunning), the sensor just keeps waiting, bound only by its owntimeout. - Omitting
external_task_idchecks the external DAG run's own state column, not a derived status of its tasks — there can be a short lag between "all its tasks look done" and the run itself flipping tosuccess.
Example:
from dag_parser.dynamic.operators import ExternalTaskSensor
wait_for_upstream_dag = ExternalTaskSensor(
task_id="wait_for_upstream_dag",
external_dag_id="daily_ingest",
external_task_id="load_final_table", # omit to wait on the whole DAG run instead
allowed_states=["success"],
failed_states=["failed", "upstream_failed"],
execution_date="{{ ds }}T00:00:00", # pin to a specific run rather than "latest"
poke_interval=60,
timeout=3600,
mode="reschedule",
)50Sensor tuning#
What it does: Three knobs common to every sensor.
Parameters:
| Param | Type | Default | Notes |
|---|---|---|---|
poke_interval | integer (seconds) | 60 | A value of 0 or negative is treated as unset and falls back to 60. |
timeout | integer (seconds) | 3600 | Cumulative across the entire wait, including every reschedule cycle. |
soft_fail | boolean | False | Timeout (and, for ExternalTaskSensor, the external target failing) becomes skipped instead of failed. |
How it works:
timeoutis cumulative from the sensor's very first check, across everyreschedule → scheduled → queued → runningcycle — it is not reset each time the sensor is re-dispatched. Design yourtimeoutaround the total acceptable wait, not a per-check budget.soft_failonly changes the timeout outcome (and, forExternalTaskSensor, the external-target-failed outcome) — a genuine poke error (bad SQL, an uncompilable regex, a missing required param) still fails the task outright, regardless ofsoft_fail.- There's no built-in minimum or backoff on
poke_interval— a very small interval on amode="poke"sensor polls its target at a fixed, un-backed-off cadence for the wholetimeoutwindow.
Example:
tight_poll = HttpSensor(
task_id="tight_poll",
endpoint="https://api.example.com/v1/status",
poke_interval=15, # checked every 15s...
timeout=600, # ...for up to 10 minutes total, cumulative
soft_fail=True, # timeout -> skipped (downstream must tolerate this via trigger_rule)
)51Poke vs. reschedule#
What it does: mode="poke" holds the worker slot for the sensor's entire wait. mode="reschedule" checks once, then releases the slot entirely until the next check is due.
Parameters:
| Sensor | Default mode |
|---|---|
HttpSensor | "poke" |
SqlSensor | "reschedule" |
TimeSensor | "reschedule" |
ExternalTaskSensor | "reschedule" |
How it works:
- The four built-in sensors do not share a common default —
HttpSensordefaults to"poke", while the other three default to"reschedule". Setmode=explicitly if you want consistent behavior across sensor types in your DAG. mode="poke"is cheap on the database (one claim, one running row, no re-dispatch churn) but ties up a live worker slot for the full wait — fine for short waits, wasteful for long ones.mode="reschedule"re-claims the task on every poke, which increments itstry_numbereach time — this is normal and expected; it doesn't consume any of the task'sretriesbudget, so don't read a largetry_numberon a reschedule-mode sensor as "it failed and retried N times."- Both modes share the exact same cumulative
timeoutsemantics (feature 50) — the mode only changes how worker resources are used while waiting, not how long the sensor is willing to wait overall.
Example:
# poke: holds the worker slot the whole time (HttpSensor's default — fine for short waits)
quick_check = HttpSensor(task_id="quick_check", endpoint="https://api.example.com/health",
poke_interval=10, timeout=120) # mode="poke" implicitly
# reschedule: releases the slot between checks (better for long waits)
long_wait = HttpSensor(task_id="long_wait", endpoint="https://api.example.com/status",
poke_interval=300, timeout=3600 * 6, mode="reschedule")52Deferrable/async waits#
What it does: A third waiting strategy, alongside poke and reschedule: a deferred task consumes zero worker-pool resources while waiting — see feature 28 for the full self.defer() walkthrough.
How it works:
- None of the four built-in sensors (features 46–49) use
defer()internally — there's nodeferrable=Truekwarg to get the zero-footprint behavior "for free" onHttpSensor/SqlSensor/TimeSensor/ExternalTaskSensor. To get fully deferred, zero-worker-footprint waiting for an HTTP- or time-based condition, write your ownBaseOperatorsubclass and callself.defer(...)yourself (feature 28). soft_failhas no equivalent for a deferred wait — a trigger that times out (per thetimeout=passed todefer()) fails the task directly; there's no "timeout as skip" option the way sensors get viasoft_fail.execute_complete(self, event=None)is where you handle the fired trigger's result — if your custom operator doesn't override it, the task simply succeeds the moment the trigger fires, without inspecting what the trigger returned.- Because the task is fully out of the worker pool while deferred, it doesn't compete for
MaxLocalTasksor any pool slot at all until the trigger fires and it re-enters as a normal scheduled task.
Example:
from dag_parser.dynamic.dag_context import BaseOperator, HttpTrigger
class DeferredHealthCheck(BaseOperator):
"""Waits for an endpoint to go healthy WITHOUT holding a worker slot
(contrast with HttpSensor, feature 46, which always holds/reschedules
through the worker pool)."""
operator_name = "DeferredHealthCheck"
def execute(self, context):
self.defer(
trigger=HttpTrigger(endpoint="https://api.example.com/health", expected_status=200),
method_name="execute_complete",
timeout=1800, # triggerer fails the task if this isn't reached
)
def execute_complete(self, event=None):
print(f"endpoint became healthy: {event}")