New to Kestra?
Use blueprints to kickstart your first workflows.
Trigger an Airflow DAG, classify its run by state and SLA duration, alert on risk, and optionally retrigger a failed run.
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.
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.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.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.log_audit_complete always runs last, printing the headline numbers regardless of which branch fired.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".auto_remediate, dry_run) between "run failed" and "retry actually triggered".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.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.
dag_id already deployed and unpaused.POST /api/v1/dags/{dag_id}/dagRuns).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.
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.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.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.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.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.AIRFLOW_USERNAME, AIRFLOW_PASSWORD, and SLACK_WEBHOOK_URL secrets.log_sla_healthy fires on a normal, fast run.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.auto_remediate: "true" and dry_run: "true" and confirm log_dry_run_remediation describes the retry without starting it.auto_remediate: "true" and dry_run: "false" and confirm retry_failed_dag_run starts a new run, visible in the Airflow UI.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.dag_id values with io.kestra.plugin.core.flow.Loop to run the same SLA gate across every critical DAG in one flow.retry_failed_dag_run, so a persistently broken DAG does not get retriggered indefinitely on every scheduled run.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.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.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.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.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.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.