Schedule icon
TriggerDagRun icon
If icon
SlackIncomingWebhook icon
Log icon
Switch icon

Airflow DAG Run SLA Gate

Trigger an Airflow DAG, classify its run by state and SLA duration, alert on risk, and optionally retrigger a failed run.

Categories
DataInfrastructure

Airflow's own UI shows a DAG's last run state if someone opens it, but nothing pages anyone when a critical DAG quietly starts taking three times as long as usual, or fails outright at 3am with nobody watching. This blueprint uses io.kestra.plugin.airflow.dags.TriggerDagRun to run a named DAG and wait for its terminal state, classifies the result by both outcome and duration, alerts on anything but a clean bill of health, and can kick off exactly one retry run when the monitored run failed outright, behind two independent gates.

How it works

  1. periodic_sla_check (io.kestra.plugin.core.trigger.Schedule) runs every hour, always forcing auto_remediate: "false" and dry_run: "true" through its own inputs: override regardless of this flow's defaults. Shipped disabled so you can validate a manual run first.
  2. trigger_monitored_dag (io.kestra.plugin.airflow.dags.TriggerDagRun, wait: true) starts dag_id and polls every 5 seconds for up to 30 minutes, returning dagRunId, state (success or failed), started, and ended once the run reaches a terminal state - or throwing if it never does.
  3. evaluate_run_state (io.kestra.plugin.core.flow.If) branches on state == "success".
    • else (UNREACHABLE): the run's own state is failed. route_remediation (io.kestra.plugin.core.flow.Switch on auto_remediate) either considers a retry ("true", gated again by dry_run through check_dry_run) or only logs that remediation is disabled ("false"). When both gates allow it, retry_failed_dag_run starts one fresh, non-blocking run of the same dag_id. alert_run_failed always fires to Slack regardless of which path ran.
    • then (run succeeded): nested evaluate_sla (io.kestra.plugin.core.flow.If) branches on (ended | timestamp) - (started | timestamp) > sla_duration_seconds.
      • then (AT_RISK): alert_sla_breach reports the measured duration against the SLA - there is no retry path here, since the run already succeeded.
      • else (HEALTHY): log_sla_healthy records a clean, on-time run.
  4. log_audit_complete always runs last, printing the headline numbers regardless of which branch fired.
  5. The errors block alerts Slack separately if the flow itself fails outright - an unreachable Airflow webserver, bad credentials, or the monitored run blowing through maxDuration without ever reaching a terminal state must not read as "no risk".

What you get

  • A direct answer from Airflow's own DAG run state and timestamps to "did this run succeed, and was it fast enough" - not inferred from a dashboard nobody is watching.
  • Outright failure (UNREACHABLE) and slow-but-successful (AT_RISK) treated as genuinely different severities, with different responses: one is eligible for an automatic retry, the other is report-only by design.
  • Remediation that only ever starts one fresh, non-blocking run of the named DAG - it never touches the original failed run's state, logs, or task instances.
  • Two independent gates (auto_remediate, dry_run) between "run failed" and "retry actually triggered".
  • A final status log every run, risk or not, so the execution history doubles as an SLA trend for that DAG.

Who it's for

  • Data platform teams running Airflow alongside Kestra who want one critical DAG watched for both outright failure and SLA creep, without building a second monitoring stack.
  • On-call engineers who want a single Slack alert naming the exact run ID and measured duration instead of opening the Airflow UI to check.
  • Anyone evaluating the Airflow plugin who wants a worked example of TriggerDagRun driving both a monitored run and a separate, non-blocking retry call, distinct from this repo's existing airflow-trigger-dag.yaml, which only triggers and waits with no classification or retry logic at all.

Why orchestrate this with Kestra

Triggering a DAG through Airflow's own UI or API answers one moment in time; it does not decide a cadence, classify the result by both state and duration, gate a retry behind two independent confirmations, or notify anyone. Kestra supplies the schedule, the HEALTHY/AT_RISK/UNREACHABLE classification as a first-class branch, an authorization gate pair before any retry run is started, and an execution history that shows exactly when a DAG started failing or slowing down.

Prerequisites

  • A reachable Airflow instance (2.x, REST API enabled) with dag_id already deployed and unpaused.
  • A user with permission to trigger DAG runs via the Airflow REST API (POST /api/v1/dags/{dag_id}/dagRuns).
  • A Slack incoming webhook for alerts.

Local testing: docker run -d --name airflow-sla-gate -p 8090:8080 -e AIRFLOW__CORE__LOAD_EXAMPLES=true -e _AIRFLOW_WWW_USER_USERNAME=airflow -e _AIRFLOW_WWW_USER_PASSWORD=<your-local-password> apache/airflow:2.9.3-python3.11 standalone docker network connect YOUR_KESTRA_NETWORK airflow-sla-gate (check existing networks first with docker network ls; only needed if the Kestra Worker runs in a separate Docker network than this container; set airflow_base_url to http://airflow-sla-gate:8080 from inside that network, or http://localhost:8090 from the host). Host port 8090 is used here specifically so it never collides with Kestra's own UI on 8080, since Airflow's webserver also defaults to 8080 internally.

Secrets

  • AIRFLOW_USERNAME / AIRFLOW_PASSWORD: basic-auth credentials used by every io.kestra.plugin.airflow.dags.TriggerDagRun task in this flow.
  • SLACK_WEBHOOK_URL: Slack incoming webhook used by alert_sla_breach, alert_run_failed, and the errors block.
  • In Kestra OSS (no Enterprise secrets backend), secrets are supplied as environment variables prefixed SECRET_, base64-encoded, and read back in flows with {{ secret('NAME') }} - for example SECRET_AIRFLOW_PASSWORD=$(echo -n '<your-local-password>' | base64). This keeps credentials out of the flow YAML but, per Kestra's own documentation, offers no encryption at rest or access control beyond the host environment; use the Enterprise secrets backend for stronger guarantees.

Inputs

  • airflow_base_url (STRING, default http://localhost:8090): Airflow base URL.
  • dag_id (STRING, default example_astronauts): DAG triggered and monitored.
  • sla_duration_seconds (INT, default 600): duration ceiling for a successful run.
  • auto_remediate (SELECT: "false", "true"; default "false"): must be "true" for a retry to even be considered.
  • dry_run (SELECT: "true", "false"; default "true"): must be explicitly "false", together with auto_remediate: "true", for a retry run to actually start.

Outputs

  • outputs.trigger_monitored_dag.dagRunId / .state / .started / .ended: the monitored run's identity and terminal state, every run.
  • outputs.retry_failed_dag_run.dagRunId: the new run's ID, present only when a retry actually fired.

Quick start

  1. Start Airflow locally with the command above (the standalone command seeds an admin user and loads example DAGs, including example_astronauts) and confirm the DAG is unpaused in the Airflow UI at http://localhost:8090.
  2. Add the AIRFLOW_USERNAME, AIRFLOW_PASSWORD, and SLACK_WEBHOOK_URL secrets.
  3. Run the flow manually with the defaults and confirm log_sla_healthy fires on a normal, fast run.
  4. Lower sla_duration_seconds to something the run will exceed (for example 1) and re-run to confirm alert_sla_breach fires (AT_RISK) with no retry option offered.
  5. Pause or break the DAG so a run fails, then re-run with auto_remediate: "true" and dry_run: "true" and confirm log_dry_run_remediation describes the retry without starting it.
  6. Re-run with auto_remediate: "true" and dry_run: "false" and confirm retry_failed_dag_run starts a new run, visible in the Airflow UI.
  7. Enable periodic_sla_check once you trust the check; it always runs in the safe auto_remediate: "false" / dry_run: "true" mode regardless of what you leave the flow's own defaults set to.

How to extend

  • Loop over several dag_id values with io.kestra.plugin.core.flow.Loop to run the same SLA gate across every critical DAG in one flow.
  • Add a maximum retry count by reading a KV-stored counter before retry_failed_dag_run, so a persistently broken DAG does not get retriggered indefinitely on every scheduled run.
  • Pass body.conf on trigger_monitored_dag with run-specific parameters if the monitored DAG expects input configuration, following the plugin's own documented default-conf shape.
  • Route the UNREACHABLE branch to a page instead of a Slack message, since an outright DAG failure is typically a stronger signal than an SLA breach on an otherwise-successful run.

Pitfalls

  • This flow triggers a new run every time it executes - it does not just inspect the last one. Unlike a passive status poll, trigger_monitored_dag starts real work in Airflow on every execution, scheduled or manual. Do not point dag_id at a DAG where an extra, unplanned run would be unsafe or expensive.
  • TriggerDagRun's Output has no duration field. The SLA check computes (ended | timestamp) - (started | timestamp) itself; there is nothing like outputs.trigger_monitored_dag.durationSeconds to read directly.
  • A maxDuration timeout throws, it does not return a third state. A run that never reaches success or failed within 30 minutes fails the task outright and is caught by the errors block, not by evaluate_run_state - that branch only ever sees the two state strings the plugin itself polls for.
  • A slow-but-successful run (AT_RISK) is never retried. route_remediation only exists on the failed-state branch; retriggering an already-successful run would not address a timing problem, so this flow deliberately offers no remediation action there.
  • The retry only starts a new run; it never inspects or clears the original failed run. retry_failed_dag_run does not mark the failed run as resolved, clear its task instances, or check whether the retry itself later succeeds - that requires a second pass of this same flow or a separate check.
  • auto_remediate and dry_run are independent gates, not a single boolean. Both must be "true"/"false" respectively for a retry to actually start; setting only one leaves the other still blocking.
  • Switch case keys "true"/"false" are quoted strings. auto_remediate renders as the literal string "true" or "false"; unquoted true:/false: YAML map keys would parse as booleans instead and would not match.
  • This blueprint is UNTESTED against a live Airflow instance. TriggerDagRun's properties (baseUrl, dagId, wait, pollFrequency, maxDuration, options) and outputs (dagId, dagRunId, state, started, ended) are taken verbatim from the plugin's source on main, and the (date | timestamp) duration pattern is copied from this repo's own dropbox-backup-to-s3.yaml - but no Docker container or Kestra engine was run to execute this flow end to end.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.