You're on Astronomer: Keep Astro and Otto on the DAG, and Let Data Workers Follow the Incident Out of It
Astro runs Airflow and Otto investigates DAG failures. Data Workers follows the incident into Kafka, Snowflake and dbt, fixes the cause behind approvals and verifies the rerun.
Your Airflow lives on Astro. Each Deployment runs Astro Runtime (3.3 on Airflow 3), Hosted by Astronomer or with Remote Execution Agents that keep tasks, code and secrets in your own infrastructure. Your dbt project runs as Cosmos task groups, Astro Observe groups DAGs and tables into data products with SLAs, and Astro alerts page the on-call. Since Sept 30, 2026, the new Astro experience is on for every Organization, Otto is available across Astro and the Astro IDE is generally available. When a DAG run fails, an Otto investigation reads the task logs, the Deployment and Observe lineage, then ranks up to three fixes for a person to apply.
Astro is where your DAGs run and Otto works on them. Data Workers is the responder that follows an incident past the DAG, into the source, the warehouse and the dbt models, and fixes the cause behind your approvals. Plenty of failed runs aren't about the DAG at all.
Key takeaways
- •Astro and Otto stay as they are. Deployments, Cosmos, Observe, alerts and Otto keep their jobs. Data Workers works through each Deployment's Airflow REST API.
- •The incident gets followed out of the DAG. When a run fails on bad input, Data Workers traces it into Kafka Connect, Snowflake and dbt and maps every downstream consumer.
- •Fixes land where the cause lives, behind approvals. Cleanups and config changes go to their owners to apply; dbt changes arrive as diffs to merge.
- •The rerun is drafted, approved and verified. Data Workers specifies the one DAG run and its
conf, a named person approves, and on Airflow 3 the owner triggers it (on Airflow 2,trigger_airflow_dagqueues it). Astro schedules, runs and records it. - •One receipt per incident. Cause, approvers, run, checks and undo in one record, across every Airflow you run.
Astro is your Airflow home, with Otto on the DAG. Data Workers is the responder that follows the incident past it.
Otto's investigations are built on what Astronomer knows best: in Astronomer's words, Otto "investigates Dag failures on Astro using proprietary Airflow, Astro, and Observe context". Its suggested fixes come in three types, a code patch, an Airflow operation (clear a task or run, or trigger a new one) and manual instructions, and a person applies them. Inside the DAG, that is the right first look.
Here is a Friday on Astro where the DAG was fine. This is an illustration, not a customer case. The product emits usage events to the Kafka topic product.usage_events. A Kafka Connect Snowflake sink lands them in RAW.PRODUCT.USAGE_EVENTS, with each record's topic, partition and offset in RECORD_METADATA. On the Astro Deployment prod-analytics, the DAG usage_billing_daily runs at 05:00. A Cosmos task group builds stg_usage_events and fct_daily_usage, and the last task, push_usage_to_stripe, sends each customer's metered usage to Stripe Billing, where usage invoices finalize at 10:00. Observe tracks a data product called "Usage billing" with a 07:00 SLA.
| Time (Fri Oct 2, ET) | System | What happens |
|---|---|---|
| Thu 23:12 | Kafka Connect | A worker restarts after running out of memory. The sink's tasks rebalance and resume from their last committed offsets, so records from four partitions are written to Snowflake a second time |
| Thu 23:58 | Kafka Connect | The same worker restarts again; another block of records lands twice. By morning RAW.PRODUCT.USAGE_EVENTS holds 3.1 million duplicate rows |
| 05:00 | Astro | The usage_billing_daily DAG run starts |
| 05:21 | Astro + Cosmos + dbt | stg_usage_events builds; its unique test on event_id fails at error severity. The task fails after its retries, and push_usage_to_stripe is marked upstream_failed. Nothing has gone to Stripe |
| 05:24 | Astro alerts + Otto | A Dag Failure alert pages the data on-call, and its Dag Trigger channel starts an Otto investigation. Observe shows "Usage billing" behind its 07:00 SLA |
| 05:26 | Data Workers | It reads the failed DAG run and its task instances through the Deployment's Airflow REST API and opens an incident |
| 05:31 | Otto | The investigation names the failing test from the task log and ranks fixes inside the DAG, such as a code patch that dedupes stg_usage_events on event_id, or clearing the task to retry. Both are reasonable for the DAG it can see |
| 05:38 | Data Workers + Snowflake + Kafka Connect | It queries the raw table: every duplicate pair shares the same topic, partition and offset, so the app sent each event once and the sink wrote it twice. It reads connector and task status from the Kafka Connect REST API: the connector is running again after two task restarts on one worker, so the worker needs a look before tonight's load |
| 05:42 | Data Workers | Blast radius from lineage: fct_daily_usage and the Stripe push for 412 customers on usage plans, whose Oct 1 usage would read about 1.3 times too high; and fct_feature_adoption, a product model that reads the raw table directly and already built at 04:00 with inflated counts |
| 05:46 | Data Workers + Slack | It posts the plan to the billing data owner and the platform on-call: a cleanup that removes rows whose partition and offset repeat, keeping the first load, with the removed rows copied to a holding table first as the undo; the worker memory fix for the platform team; a rerun of usage_billing_daily for Oct 1; a rebuild of fct_feature_adoption; and a staging guard that dedupes on partition and offset, as a diff for the dbt owner |
| 06:05 | Spellbook + Snowflake + Kafka Connect | The billing data owner approves the plan in Spellbook and applies the cleanup in Snowflake. The platform on-call raises the worker's memory and restarts the worker themselves |
| 06:16 | Data Workers + Astro | It confirms no partition and offset repeats in the raw table and marks the rerun ready; the billing data owner triggers the usage_billing_daily run it drafted, with {"usage_date": "2026-10-01"}, and Airflow runs it |
| 06:41 | Astro + Cosmos + dbt | stg_usage_events and fct_daily_usage build; every test passes |
| 06:44 | Astro + Stripe | push_usage_to_stripe sends each customer's Oct 1 usage once |
| 06:49 | Data Workers | It verifies: raw row count equals distinct partition and offset pairs, usage per customer sits inside its trailing two-week range, and the usage total the billing owner confirms in Stripe equals fct_daily_usage. The receipt lands in Spellbook, and the on-call links it from the Otto investigation |
| 06:55 | Astro Observe | The "Usage billing" data product finishes inside its 07:00 SLA |
| 10:00 | Stripe | Usage invoices finalize on events counted once |
| 11:20 | GitHub | The dbt owner merges the staging guard; fct_feature_adoption rebuilds through its own DAG after approval |

Astro did its job: the test failed loudly, the push never ran, Observe flagged the SLA and Otto explained the failure. The cause sat two systems upstream, and part of the damage had already reached a model the failed DAG doesn't touch. This domain runs at L2 propose, so every change waited for its owner.
| Job | What Astro does | What Data Workers does |
|---|---|---|
| Running the work | Hosts the Deployment, schedules the DAG, runs dbt through Cosmos and retries failed tasks | Reads DAG runs and task instances through the Airflow REST API and ties each task to the tables it builds |
| Raising the alarm | Observe SLAs on data products and Astro alerts to the on-call | Opens an incident with the cause traced across systems, including runs that fail on bad input |
| Investigating | Otto reads task logs, the Deployment and Observe lineage, and ranks up to three DAG fixes | Follows the failure into Kafka Connect, Snowflake and dbt, and maps every consumer it reaches, inside and outside the DAG |
| Fixing it | A person applies Otto's code patch, Airflow operation or manual steps | Proposes each fix where the cause lives: a cleanup for the data owner, a config change for the platform team, a dbt diff to merge |
| Running it again | Runs the new DAG run and keeps its record | Drafts the approved run (DAG, conf, order) and confirms the inputs check clean before the owner triggers it on Airflow 3, with the undo written into the plan |
| Proving it | Marks the run green and the SLA met | Checks the rows against the source and the push against the model, then leaves a receipt: cause, approvers, run, checks, undo |
Why doesn't Astro just do this itself?
Because Astronomer built the best place to run Airflow, and scoped its agent to match. Otto investigations draw on Airflow, Astro and Observe context, and their fixes apply to a DAG or a Deployment: a patch to one file, a clear or a trigger, written steps, applied by people. In an agent session Otto can also query your warehouse and trace lineage. The hosted Astro MCP server, in Labs, gives MCP clients read access to the Astro control plane under the user's own permissions. Each choice fits a vendor running many companies' Airflow.
Removing three million rows from a Snowflake table, changing a Kafka Connect worker and deciding whether 412 customers get billed today is a different job with a different liability. It needs context about systems Astro doesn't run, a blast radius past the DAG, the right owner for each change, an undo for every step and verification against the source. That cross-system operations product, with approvals and rollback built in, is what Data Workers, the agentic data platform, runs. Otto stays the DAG expert; Data Workers takes the incident when it leaves the DAG.
Every tool owns a slice. Data Workers covers the whole lifecycle
Astro owns its slice deeply: running Airflow, from Deployments and Cosmos to Observe and Otto. Each point tool adds another console, contract and handoff. Data Workers covers the whole lifecycle with one context, one approval flow and one audit trail, building on the Deployments already running your schedule.

| Stage | Data Workers | Astro | Why we scored it this way |
|---|---|---|---|
| Catalog & Context | 9 | 3 | Observe groups DAGs, tasks and tables into data products with lineage from OpenLineage. Data Workers keeps one governed context graph of tables, models, owners and lineage across every system. |
| Analytics & Insights | 8 | 1 | Not Astro's job: dashboards report on deployments, runs and SLAs, not on the business numbers the DAGs build. Data Workers answers data questions from governed definitions with lineage behind every number. |
| Data Quality | 8 | 3 | dbt tests run inside Cosmos task groups and checks run as tasks you write. Data Workers writes, runs and repairs quality checks on the tables the DAGs build. |
| Observability & Incidents | 8.5 | 7 | Observe SLAs and AI log summaries, Astro alerts, and Otto investigations that rank DAG fixes. Data Workers traces the cause into the source and the warehouse, proposes the fix behind an approval and verifies the data. |
| Pipelines & Ingestion | 8.5 | 9 | Astro's home stage: managed Airflow deployments, Hosted or Remote Execution, the Astro IDE, Cosmos for dbt and Otto for authoring. Data Workers drafts approved reruns, checks them and keeps what they load right. |
| Schema & Migration | 8 | 2 | Not Astro's job: a task sees success or failure, not a changed record, a re-delivered offset or a new schema. Data Workers catches schema and shape changes and assesses blast radius. |
| Governance & Access | 8.5 | 4 | Workspace and deployment roles control who can deploy, trigger and clear. Data Workers routes every data change to a named approver with a receipt. |
| Security & Privacy | 8 | 3 | Secrets, connections, IP allowlists and Remote Execution keep Airflow itself secure. Data Workers' pull request review flags new columns whose names or annotations look sensitive before the DAGs move them. |
| Cost / FinOps | 8 | 4 | Observe analyzes Snowflake credit use per task, and Astro reports its own AI spend. Data Workers traces Snowflake spend to the dbt model behind it and drafts the fix for its owner. |
| MLOps & Models | 7.5 | 4 | Training and scoring pipelines run as DAGs, so Astro schedules the steps. Data Workers keeps the data under your models healthy. |
Running Airflow outside Astro too? You're on Airflow covers self-managed Airflow, Managed Airflow and MWAA, and you're on Google Cloud Composer covers Google's managed service. The Kafka side of this story is in you're on Apache Kafka, the dbt side in you're on dbt.
How Astro and Data Workers work together
Engineers ask from their coding agent, with Otto in the Astro IDE or CLI; approvers decide in Spellbook Data Catalog (in preview). Underneath, Data Context Wizard keeps one governed context graph across Astro, Kafka, Snowflake, dbt and Stripe, the Data-Agents Swarm does the work, and the Autonomous Data-Conductor runs each incident end to end: detect, diagnose, fix, review, verify, remember.

What Data Workers reads and triggers on Astro. Every Astro Deployment exposes the Airflow REST API, and Data Workers' native Airflow connector uses it with an Astro API token, detecting /api/v2 on Airflow 3 or /api/v1 on Airflow 2. It reads DAG runs and their state and the task instances in each run. Lineage comes from the dbt manifest and the team's context-graph notes. On Astro Runtime for Airflow 3, the owner triggers each approved run from the plan Data Workers drafted, and Data Workers verifies the result. On an Airflow 2 Deployment, the one thing it writes is a new DAG run, with an optional conf payload, queued through trigger_airflow_dag after a named approval. Clearing tasks, Airflow backfills, pausing DAGs, deploys and Deployment settings stay with your team, by design. Use a Deployment-scoped API token with the narrowest role that covers what Data Workers does there, or pin the connection read-only and every write is refused. The wiring is in the Airflow integration guide.
The rest of the stack. Snowflake and dbt are native. Kafka Connect is native too: Data Workers reads connector and task status over its REST API, and restarts and worker changes stay with the platform team. Stripe connects over its API, read only: Data Workers checks what the DAG sent and never pushes to Stripe. Otto, Observe and the Astro API connect over their APIs or MCP servers today.
Backfills, precisely. When the fix is a reload, Data Workers drafts it as a rerun for the owner to trigger on Airflow 3, reads its status, and remediate re-checks the quality assertions and escalates any failure to a person. The undo is written into the plan before the run, and the owner runs it if it's ever needed. Cleanups such as the dedupe above are proposed for the data owner to approve and apply; Data Workers deletes nothing in the warehouse.
Setup over MCP today. Every Data Workers agent is an MCP server. Per the client setup guide, clone the repo and add one start-agent.sh entry per agent to your client.
# Example: Data Workers agents in Claude Code, from a clone of the open-source repo
claude mcp add --scope user dw-incidents -- "$(pwd)/start-agent.sh" dw-incidents
claude mcp add --scope user dw-catalog -- "$(pwd)/start-agent.sh" dw-context-catalog
claude mcp add --scope user dw-quality -- "$(pwd)/start-agent.sh" dw-quality
claude mcp add --scope user dw-connectors -- "$(pwd)/start-agent.sh" dw-connectorsList each agent's tools with your client's own command, such as /mcp. In this story: diagnose_incident explains the failure; trace_cross_platform_lineage and blast_radius_analysis find fct_feature_adoption and the Stripe push; run_quality_check and get_quality_score confirm the raw table is clean; send_slack_alert tells the channel. On this Airflow 3 Deployment the owner triggers the approved run; on Airflow 2, trigger_airflow_dag queues it. Astronomer's Astro MCP server and astro-airflow-mcp can sit in the same client.
Where things run. The agents run in your infrastructure on every tier and hold the Astro, Snowflake and model credentials. Your data stays in your systems; the hosted Conductor sees workflow metadata only. The remote endpoint takes an API key or OAuth tokens from your identity provider, such as Okta or Entra ID, verified through JWKS.
One incident, L0 to L4, set per domain:

- •L0 manual. The on-call applies Otto's dedupe patch, gets the run green and finds the adoption numbers wrong a week later.
- •L1 observe. The incident opens at 05:26 with cause and blast radius. Nothing changes.
- •L2 propose. Nothing runs until a named owner approves; an unanswered request expires and escalates, never auto-grants.
- •L3 act reversibly. For a class with a clean record, such as rerunning one date once its inputs check clean, the run goes to its owner as soon as the checks pass; on an Airflow 2 Deployment, Data Workers queues it.
- •L4 autonomous. A scoped domain handles that class end to end; people read receipts.
What changes for your team
Astro keeps running Airflow, your DAG authors keep authoring and Otto keeps helping them. What changes is the dawn bridge call: a failed run whose cause sits upstream arrives as one incident with the cause, the blast radius, a plan for each owner and a rerun waiting for one approval.

The Orchestration agent and Incident Debugging agent posts go deeper. For practice notes, see automating Airflow task failure analysis and reducing data on-call burden; Data Workers for Airflow has the overview.
Keep Astro, or consolidate?
Keep Astro if you love it; Data Workers works with it from day one. Many teams consolidate once Data Workers runs that slice too.
For most Astro teams the answer is keep it: Astronomer runs Airflow so your team doesn't have to. What teams consolidate is the scaffolding around it: dedupe and row-count tasks copied into every DAG, and runbooks that say "check the sink, clear the task, check the dashboard". Those move into one loop that ends with a verified Airflow run and a receipt. Weighing a build on a coding agent plus Astronomer's MCP servers? Read build it ourselves with Claude Code and MCP servers: wiring the servers is the easy part; approvals, blast radius, undo and receipts across every system are the product.
The case for your CFO
The outcome: the numbers your Astro DAGs build, including the ones that bill customers, arrive right and on time. When bad input breaks a run, the cause is found where it lives, fixed by its owner behind an approval and verified before invoices go out.
The risk story is plain. Agents decide nothing you haven't delegated, and autonomy is set per domain from L0 manual to L4 autonomous. At L2 a named person approves every rerun and fix, and no agent can promote its own work. On Astro, reruns, data cleanups and config changes are run by their owners from the approved plan, and Data Workers verifies each one. Every incident carries a receipt with the cause, approvers, run, checks and undo, and an org-wide stop halts all autonomous dispatch. Zero migration: Astro, your DAGs, Cosmos, Observe, Otto, Kafka and Snowflake stay as they are.
Why now: every tool in the stack now ships its own agent, each fixing its own slice; someone has to own the incident across them. The first win is L1 on one revenue-critical data product, such as usage billing. What stays the same: your Deployments, repository, review process, Astro roles and on-call rota. For the numbers, see the ROI of agentic data operations.
The sentence to repeat upstairs: "Astro and Otto keep our DAGs running; Data Workers follows every incident out of the DAG and fixes the cause behind our approvals."
Getting started
Start with a pilot. Pick one Astro Deployment and one data product that matters, such as usage billing, connect Data Workers read-only to that Deployment's Airflow REST API, the warehouse, the dbt project and the source that feeds it, and run at L1 so every failed or suspect run arrives with its cause and blast radius. Then turn on L2 for that domain, so every fix and rerun arrives as an approved plan for its owner. Plans are on the pricing page, and the pilot is credited in full against the first year. Data leaders can read the Airflow guide for data leaders.
FAQ
We already use Otto. Do we need Data Workers? Keep Otto for authoring, reviewing and investigating DAGs on Astro; it has context about your Deployments no general tool has. Data Workers works beside it and takes the incidents whose cause or impact sits outside the DAG: a source, a sink connector, a warehouse table, a model another DAG builds.
What does Data Workers read from Astro, and what does it change? It reads DAG runs and task instances through each Deployment's Airflow REST API, not the DAG list. On Airflow 3 it changes nothing there: the owner triggers each approved run from Data Workers' plan, and Data Workers verifies it. On an Airflow 2 Deployment it queues a new DAG run after a named approval. Deploys, clears, Airflow backfills, pauses and Deployment settings stay with your team.
Does it work with Remote Execution, Astro Hybrid, or Airflow outside Astro? Yes. Data Workers talks to each Deployment's Airflow REST API and its agents run in your infrastructure, so the pattern is the same wherever tasks run. Self-managed Airflow, Managed Airflow and MWAA join the same approval flow and audit trail.
Will Data Workers delete rows in our warehouse, or restart our connectors? No. Cleanups like the dedupe above are proposed as SQL with the affected rows copied aside first, and the data owner approves and applies them. Kafka Connect restarts and worker changes stay with the platform team. After a reload, Data Workers re-checks the data and escalates any failure to a person.
Will it change our DAG or dbt code? It proposes the change as a diff for the owner to merge, and opens a pull request only when your team turns on the GitHub pull-request target.
Sources
- •Astronomer, Astro release notes (Sept 30, 2026: new Astro experience on for all Organizations, Otto available across Astro, Astro IDE GA, Otto investigations list up to three ranked suggested fixes, Labs, with Code patch, Airflow and Manual fix types, Astro AI cost breakdown; Sept 15, 2026: Ask Otto in the Astro IDE; Astro CLI 1.45.0, Astro Runtime 3.3-7, Remote Execution Agent 1.8.6), https://www.astronomer.io/docs/astro/release-notes (checked Oct 3, 2026)
- •Astronomer, Otto investigations (Labs; "investigates Dag failures on Astro using proprietary Airflow, Astro, and Observe context"; users copy and apply fixes; investigation API; Astro alerts can start investigations), https://www.astronomer.io/docs/astro/otto/otto-investigate (checked Oct 3, 2026)
- •Astronomer, Otto overview (Labs label; exploration: query your warehouse, trace lineage; Astro CLI and Astro IDE access), https://www.astronomer.io/docs/astro/otto (checked Oct 3, 2026)
- •Astronomer, Astro alerts (Dag Failure alerts; Dag Trigger notification channel calls the Otto investigation API), https://www.astronomer.io/docs/astro/alerts (checked Oct 3, 2026)
- •Astronomer, Astro MCP server (Labs; hosted; read access to the Astro control plane; signs in with Astro credentials), https://www.astronomer.io/docs/astro/astro-mcp-server (checked Oct 3, 2026)
- •Astronomer, Astro architecture (Hosted and Remote execution) and Astro Hybrid overview, https://www.astronomer.io/docs/astro/astro-architecture and https://www.astronomer.io/docs/astro/hybrid-overview (checked Oct 3, 2026)
- •Astronomer, Astro Observe general availability (Feb 13, 2025; SLAs for data products, lineage, AI log summaries, Snowflake cost analysis), https://www.astronomer.io/blog/observe-ga/ (checked Oct 3, 2026)
- •Astronomer, Observe and OpenLineage (data product lineage graph), https://www.astronomer.io/docs/astro/observe-openlineage (checked Oct 3, 2026)
- •Astronomer Cosmos releases (v1.15.1, Aug 4, 2026), https://github.com/astronomer/astronomer-cosmos/releases (checked Oct 3, 2026)
- •Astronomer, astronomer/agents (astro-airflow-mcp; Apache 2.0), https://github.com/astronomer/agents (checked Oct 3, 2026)
- •Data Workers open-source repository (tool registrations:
trigger_airflow_dag,send_slack_alertin dw-connectors; dw-incidents, dw-context-catalog, dw-quality), https://github.com/DataWorkersProject/dataworkers-claw-community (checked Oct 3, 2026) - •Data Workers client setup guide, https://dataworkers.io/opensource-docs/client-setup/ (checked Oct 3, 2026)