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:

ParamTypeDefaultNotes
endpointstring (URL)required
connection_idstringnoneIf set, resolves auth (bearer token/basic auth) from a stored HTTP connection.
headersdict{}Merged on top of anything resolved from connection_id, or used on its own.
response_check_regexregex stringnoneUses Go's RE2 regex syntax (no lookahead/lookbehind/backreferences).
poke_intervalinteger (seconds)60
timeoutinteger (seconds)3600Cumulative across the whole wait (see feature 49).
modestring"poke"See feature 50 — HttpSensor's default differs from the other three sensors.
soft_failbooleanFalseTimeout becomes skipped instead of failed.

How it works:

  • Set connection_id to reuse stored credentials for the target endpoint; add headers= 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_interval until timeout, rather than failing fast.
  • response_check_regex only 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 — set mode="reschedule" explicitly if you want it to free its worker slot between checks.

Example:

python
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:

ParamTypeDefaultNotes
sqlstringrequired
connection_idstringrequiredThe external database connection to run the query against.
poke_intervalinteger (seconds)60
timeoutinteger (seconds)3600Cumulative (feature 49).
modestring"reschedule"

How it works:

  • connection_id determines 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:

python
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:

ParamTypeDefaultNotes
target_timestring, "HH:MM" or "HH:MM:SS"requiredEvaluated in the DAG's own timezone (feature 3).
poke_intervalinteger (seconds)60
timeoutinteger (seconds)3600Cumulative (feature 49).
modestring"reschedule"

How it works:

  • target_time is evaluated against the DAG's own timezone setting (feature 3) — a DAG declared timezone="America/New_York" with target_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_time on 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:

python
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:

ParamTypeDefaultNotes
external_dag_idstringrequired
external_task_idstringnoneOmit to wait on the whole DAG run's own state instead of one task.
allowed_stateslist of strings["success"]Case-insensitive match.
failed_stateslist of strings["failed"]Case-insensitive match; checked first, ends the wait as a failure.
execution_datestring (ISO or templated)latest runPin 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. Pass execution_date explicitly if you need a specific, matched run.
  • If the external state matches neither allowed_states nor failed_states (e.g. it's still running), the sensor just keeps waiting, bound only by its own timeout.
  • Omitting external_task_id checks 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 to success.

Example:

python
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:

ParamTypeDefaultNotes
poke_intervalinteger (seconds)60A value of 0 or negative is treated as unset and falls back to 60.
timeoutinteger (seconds)3600Cumulative across the entire wait, including every reschedule cycle.
soft_failbooleanFalseTimeout (and, for ExternalTaskSensor, the external target failing) becomes skipped instead of failed.

How it works:

  • timeout is cumulative from the sensor's very first check, across every reschedule → scheduled → queued → running cycle — it is not reset each time the sensor is re-dispatched. Design your timeout around the total acceptable wait, not a per-check budget.
  • soft_fail only changes the timeout outcome (and, for ExternalTaskSensor, the external-target-failed outcome) — a genuine poke error (bad SQL, an uncompilable regex, a missing required param) still fails the task outright, regardless of soft_fail.
  • There's no built-in minimum or backoff on poke_interval — a very small interval on a mode="poke" sensor polls its target at a fixed, un-backed-off cadence for the whole timeout window.

Example:

python
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:

SensorDefault mode
HttpSensor"poke"
SqlSensor"reschedule"
TimeSensor"reschedule"
ExternalTaskSensor"reschedule"

How it works:

  • The four built-in sensors do not share a common default — HttpSensor defaults to "poke", while the other three default to "reschedule". Set mode= 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 its try_number each time — this is normal and expected; it doesn't consume any of the task's retries budget, so don't read a large try_number on a reschedule-mode sensor as "it failed and retried N times."
  • Both modes share the exact same cumulative timeout semantics (feature 50) — the mode only changes how worker resources are used while waiting, not how long the sensor is willing to wait overall.

Example:

python
# 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 no deferrable=True kwarg to get the zero-footprint behavior "for free" on HttpSensor/SqlSensor/TimeSensor/ExternalTaskSensor. To get fully deferred, zero-worker-footprint waiting for an HTTP- or time-based condition, write your own BaseOperator subclass and call self.defer(...) yourself (feature 28).
  • soft_fail has no equivalent for a deferred wait — a trigger that times out (per the timeout= passed to defer()) fails the task directly; there's no "timeout as skip" option the way sensors get via soft_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 MaxLocalTasks or any pool slot at all until the trigger fires and it re-enters as a normal scheduled task.

Example:

python
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}")