Product
Product11 min readBy The Data Workers Team

You're on Kestra: Your Team Declares and Runs Event-Driven Flows. Data Workers Does the Operations Work When a Flow's Output Goes Wrong

Kestra runs declarative YAML flows across languages, triggers and plugins. Data Workers catches the execution that succeeded with the wrong output, finds the cause and hands the owner an approved fix and run plan.

Your team writes every workflow as a YAML flow: an id, a namespace, tasks and triggers, each task a plugin from a catalog of more than 2,100. A Schedule trigger runs the nightly sync, a Realtime trigger starts an execution the moment a message lands on Pub/Sub or Kafka, and a Flow trigger chains the next flow. Shared state sits in the namespace KV store, every saved flow becomes a new revision, and replay re-runs an execution from the task that went wrong. Kestra 2.0 shipped on September 7, 2026 with a rebuilt engine, flows as agent tools and a persistent AI Copilot, while the 1.3 and 1.0 lines keep getting patches.

The hardest Kestra incidents end in SUCCESS. Every task ran, every execution turned green, and the numbers in the warehouse are still wrong because an input was wrong. Data Workers watches the tables your flows write, traces a bad number back to the execution and the input behind it, and hands the owner the fix and a run plan to approve.

Key takeaways

  • •Kestra keeps its job. Flows, triggers, namespaces, the KV store and revisions stay where your team declared them. Data Workers works on what the executions produce.
  • •Green executions get checked too. Volume and null checks on the tables your flows write, including your own audit tables, turn a successful execution with wrong output into a diagnosed incident the same night.
  • •Connected over the Kestra API and MCP today. Your assistants reach Kestra executions and flows as MCP tools next to the Data Workers agents, and Data Workers checks the tables the flows write.
  • •Every fix goes through a named person. The owner approves the flow change, the trigger hold and the run plan, and carries them out in Kestra; each step leaves a receipt.
  • •Autonomy is set per domain. Start at L1 observe, move to L2 propose, and open L3 act reversibly for narrow classes once the record supports it.

Kestra is the declared plan for every flow. Data Workers is the operations desk for what they produce.

Kestra starts the execution when the event arrives, runs each task in the language it was written in, retries what you told it to retry and records every state. The operations work around a bad output is a different job: notice that a green night produced wrong numbers, connect them to another team's change, size the damage, hold anything that would ship the error to customers, fix the cause, rerun the right flows in order and prove the result. Here is one night at an EV charging network. This is an illustration, not a customer case.

TimeSystemWhat happens
Tue 16:10Pricing serviceThe pricing team ships tariff API v2 with October tariffs effective Wednesday 00:00 (DC fast charging moves from 0.49 to 0.56 per kWh). The v1 endpoint stays up for older clients and serves a frozen snapshot of the September tariffs with HTTP 200
Wed 00:05Kestrapricing.tariff_sync runs on its Schedule trigger, calls /v1/tariffs/current, gets 200 and writes the September tariffs to the KV key tariffs_current. The execution is SUCCESS
00:06Pub/Sub + Kestra + BigQueryThe charge point system publishes a SessionStopped event to Google Cloud Pub/Sub for every finished session. A Realtime trigger starts one execution of energy.rate_session per event; its Python task reads {{ kv('tariffs_current') }}, prices the session and writes it to raw.rated_sessions. Every execution succeeds
02:00Kestra + dbtenergy.revenue_marts runs dbt build: fct_session_revenue and agg_site_revenue_daily rebuild, and every dbt test passes, because 0.49 is a valid price
02:20Kestra + BigQueryThe finance team's finance.rate_audit flow joins the night's sessions to the pricing database through a BigQuery federated query and writes every session whose applied rate differs from the tariff in effect to audit.rate_mismatches. It is a report, not a gate, so it succeeds
02:30Data Workers + BigQueryThe volume check on audit.rate_mismatches reads 1,946 rows against a norm of zero. Data Workers opens an incident
02:36Data Workers + Kestra API + PostgreSQLThe on-call's assistant reads the night's executions over Kestra's MCP tools and hands them to Data Workers: tariff_sync succeeded at 00:05 and logged a call to /v1. The tariff table in the pricing service's PostgreSQL lists October rates from 00:00. Blast radius: two revenue models, the Superset "Network revenue" dashboard (read over Superset's API) and billing.fleet_statements, a Schedule-triggered flow that emails statements to fleet customers at 07:30
02:40Opsgenie + Teams + emailData Workers raises an Opsgenie alert with the diagnosis, posts the finding card to the data team's Teams channel and emails the approval request to the owner of the pricing flows
06:15Spellbook + KestraThe owner reviews three proposals in Spellbook: hold the 07:30 statements; a YAML diff that points tariff_sync at /v2 and adds an Assert task that fails the execution if effective_from is older than the current tariff period; and a run plan (tariff sync, the team's energy.rerate_day flow for Wednesday, the revenue marts, then the rate audit). She approves, disables the statements trigger in Kestra and saves the diff as revision 9 of tariff_sync
06:22Claude Code + Kestra MCPFrom her Claude Code session she starts the approved run plan through Kestra's MCP server, where each of those flows is a tool. energy.rerate_day re-prices the 2,203 sessions since midnight and merges them into raw.rated_sessions on session_id. Kestra records SUCCESS for each by 07:02
07:05Data Workers + BigQueryData Workers verifies: the rerun audit wrote zero mismatches, session counts in raw.rated_sessions match the night's events with no duplicate session_id, and both marts rebuilt
07:06OpsgenieData Workers closes the alert with the receipt: cause, approvals, runs, checks passed and the undo
07:10 to 07:30KestraThe owner re-enables the statements trigger. At 07:30 fleet customers get statements on the right rates
08:00SupersetThe operations review opens the revenue dashboard on verified numbers
Incident timeline across the stack: what Kestra, your team and Data Workers each do, step by step

Every part of Kestra did its job. The HTTP task got a 200, the KV store held what it was given, and each Realtime execution priced its session exactly as written. The fault lived in another team's service, as a price that was valid on Tuesday and wrong on Wednesday. Catching it takes knowledge Kestra was never meant to hold: which tables carry the tariff, which check says it is wrong, and which customer-facing flow reads them next.

JobWhat Kestra doesWhat Data Workers does
The runStarts executions on Schedule, Realtime and Flow triggers; runs Python, SQL, shell and dbt tasks; retries and records statesReads executions and their task logs over the Kestra API
The signalExecution states, alerts and, in Enterprise, Cases with severity, assignees and SLAsChecks the data itself: volume and nulls on the tables flows write, so a green execution with wrong output still raises an incident
The diagnosisAI Copilot in Ask mode explains a failed execution; revisions show what changed in a flowJoins the executions, the input, the source and lineage into one cause, with the blast radius across models, dashboards and downstream flows
The fixRuns whatever YAML the team saves; Copilot Edit drafts flow changes and asks before applying themProposes the flow diff and the hold to the owner, who approves and applies them in Kestra
The rerunRuns the execution it is given, from the UI, the API or an MCP tool; replay restarts from a chosen taskProposes the run plan in order; the owner starts it in Kestra, and Data Workers reads each run's status
The proofKeeps revisions, execution history and logsRe-checks the tables and writes a receipt: what changed, who approved it, how it was checked, how to undo it

Why doesn't Kestra just do this itself?

Because Kestra is built to run any workflow, in any language, against any system. Its unit is the execution: which tasks ran, with which inputs and outputs, on which revision. It does not know that 0.49 stopped being the right price at midnight, that raw.rated_sessions feeds a revenue model, or that a statements flow reads that model at 07:30.

Kestra's AI features follow the same scope, and they are careful work. The AI Copilot drafts and edits flow YAML, Edit mode asks for confirmation before applying a change, and Plan mode needs approval for each step. AI Agent tasks let a model choose which tasks or flows to call at runtime. Kestra 2.0 turns any flow with an McpToolTrigger into a tool on a tenant-scoped MCP server, private by default, and in Enterprise puts Cases in front of sensitive steps so "the agent is told the work is pending rather than done."

Owning whether a flow's output is right across a pricing service, a warehouse, dbt, BI and a customer email is a different product: a context graph of every table and downstream flow, blast-radius scoping, approvals that name a person, a recorded undo and receipts an auditor can read. That is Data Workers. For how that holds up in production, read is it safe to let AI agents change production data.

Every tool owns a slice. Data Workers covers the whole lifecycle

Kestra owns declaring and running event-driven flows. 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, and builds on the Kestra flows already there.

Spider chart of ten jobs a data team does: Data Workers covers the whole list, Kestra goes deep on its own area
StageData WorkersKestraWhy we scored it this way
Catalog & Context92Not Kestra's job: it knows flows, namespaces, executions and KV pairs, not what a table means or who owns it. Data Workers keeps one governed context graph of tables, models, lineage and owners.
Analytics & Insights81Not Kestra's job: dashboards show execution health, not business numbers. Data Workers answers data questions from governed definitions with lineage behind every number.
Data Quality83Flows can run any test tool as a task and stop on an Assert. Data Workers runs volume and null checks on the tables flows write, on every run.
Observability & Incidents8.55Execution states, alerts and, in Enterprise, Cases with SLAs track a broken execution. Data Workers diagnoses the data incident across systems, proposes the fix and verifies it.
Pipelines & Ingestion8.59Kestra's home stage: YAML flows in any language, 2,100+ plugins, Realtime, Flow and Schedule triggers, replay and backfills. Data Workers plans the reruns for the owner to start.
Schema & Migration82Not Kestra's job: a flow writes whatever schema its tasks emit. Data Workers catches schema changes and assesses blast radius before downstream models break.
Governance & Access8.54Enterprise RBAC, SSO, audit logs and MCP server permissions govern who can do what in Kestra. Data Workers routes every data change to a named approver.
Security & Privacy83Secrets, worker isolation and workers in your own network protect the runs; the data a flow writes is classified elsewhere. Data Workers leaves a receipt on every data change.
Cost / FinOps82Not Kestra's job: concurrency limits shape load, while warehouse spend sits in the warehouse. Data Workers sums BigQuery spend from the Jobs API and traces Snowflake credits to the dbt model behind them.
MLOps & Models7.54AI Agent tasks and RAG flows orchestrate model work as flows. Data Workers keeps the data under those models healthy.

How Kestra and Data Workers work together

How Data Workers fits with Kestra: your coding agent on top, Data Workers in the middle, your estate underneath

Data Workers connects to Kestra over its API or MCP server today. Kestra is API-first ("every action available in the UI is also available via REST API"), so with an API token on Enterprise or Cloud, or basic auth on open source, Data Workers reads executions, task states and logs. Flows, triggers, namespaces, KV pairs, secrets, roles and runs stay with your team: when a fix needs a flow change, Data Workers proposes the YAML as a diff, and the owner saves it as a new revision or merges it through the team's Git workflow, then starts the approved runs. The warehouse side runs natively: BigQuery, PostgreSQL and dbt in this incident, plus Snowflake, Databricks and 50+ connectors in all.

Engineers can run Kestra's MCP server and the Data Workers agents side by side: Kestra's server exposes the flows you choose as tools, and Data Workers answers "what did last night's executions do to the data, and what's the fix". The Data Workers side follows the documented client setup: clone the open-source repo and add each agent's start-agent.sh entry to your client. The Kestra side comes from the Connect tab of your MCP server in the Kestra UI.

# Example: Kestra's MCP server plus Data Workers agents in Claude Code
# Kestra MCP server: copy the URL and command from Tenant > MCP Servers > Connect
claude mcp add kestra https://kestra.example.com/<server-url-from-connect-tab> \
  --transport http \
  --header "Authorization: Basic $(echo -n 'username:password' | base64)"

# Data Workers agents, 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-quality -- "$(pwd)/start-agent.sh" dw-quality
claude mcp add --scope user dw-catalog -- "$(pwd)/start-agent.sh" dw-context-catalog

List the tools with your client's own command (/mcp in Claude Code). In this incident: run_quality_check (dw-quality) runs the volume checks; monitor_metrics (dw-incidents) tracks mismatch rows and session counts; diagnose_incident names the cause; blast_radius_analysis and trace_cross_platform_lineage (dw-context-catalog) map what the stale tariff reached; remediate re-checks the quality assertions after the reruns and escalates any failure to a person; get_incident_history shows whether the pricing feed has broken before. Keep the Kestra server private, and expose only the flows you would let an assistant start.

In production the agents run in your infrastructure and hold the warehouse credentials and model key, the same way Kestra 2.0 lets you place workers in your own network. Your data stays in your systems; the hosted Conductor sees workflow metadata only. More in where does our data go.

One incident, L0 to L4, set per domain:

The autonomy ladder: L0 manual, L1 observe, L2 propose, L3 act reversibly, L4 autonomous
  • •L0 manual. A fleet customer emails at 09:30 asking why the new tariff is missing from their statement. An engineer finds every execution green and traces the stale KV value by lunch.
  • •L1 observe. Data Workers flags the mismatches at 02:30 with the cause and blast radius. Nothing changes; the morning starts from the answer.
  • •L2 propose. Data Workers proposes the hold, the flow diff and the run plan. Nothing reaches production until the named owner approves; an unanswered request expires and escalates, never auto-grants.
  • •L3 act reversibly. For a class with a clean record, Data Workers carries out reversible steps inside the domain you open, verifies them and sends any failed check to a person. Flow changes, trigger holds and runs in Kestra stay with the owner.
  • •L4 autonomous. For a scoped domain, Data Workers watches the audit table and session volumes from the first priced sessions, so the stale input is flagged soon after midnight and the owner's fix is waiting before anything ships.

What changes for your team

Six jobs that run on autopilot with Data Workers next to Kestra, with a concrete example of each
  • •On-call starts from a cause. The engineer who picks up the alert finds the executions, the input that went wrong, the damage, the proposed fix and the run plan in one place.
  • •Customer-facing flows get a guard. When a statements, invoicing or notification flow reads a table that just went wrong, the hold is proposed before it runs.
  • •Platform and data teams share one record. Both see the same incident and receipt in Spellbook Data Catalog (in preview), linked from the Opsgenie alert.

Keep Kestra, or consolidate?

Keep Kestra if you love it; Data Workers works with it from day one. Many teams consolidate once Data Workers runs that slice too.

For most Kestra teams the answer is keep it. What teams consolidate is the tooling around the outputs: a separate observability tool, Assert tasks copied into every flow, and runbooks that say "replay and check the dashboard". Many estates also run a second orchestrator; Data Workers connects natively to Airflow, Dagster, Prefect, Azure Data Factory and Managed Service for Apache Airflow, so one incident record spans them. See you're on Airflow, you're on Prefect and you're on Dagster.

If you are weighing building this yourself on Kestra's MCP server, an AI Agent task and a coding agent, read build it ourselves with Claude Code and MCP servers. Exposing flows as tools is the easy part; the context graph, approvals, undo and receipts are the work.

The case for your CFO

The outcome. When a Kestra flow produces wrong revenue, billing or customer numbers, it is caught the same night and corrected before customers see it, with a record of what went wrong and how it was checked.

The risk story. At L0 and L1, agents only read. At L2 they propose and a named person approves; an unanswered request expires and escalates, never auto-grants. At L3 they act on reversible changes inside the domains you open; L4 is a later choice per domain. Flow changes, trigger holds and reruns stay with the owner in Kestra. No agent can promote its own work, and an org-wide stop halts all autonomous dispatch. Every change carries a receipt: what changed, who approved it, the blast radius and how to undo it.

Why now. Kestra 2.0 makes every flow callable by an AI agent. More flows, started by more callers, means more executions that succeed with the wrong output; someone has to own whether the result is right.

The first win. L1 on the flows that feed billing and customer statements: their tables get checked after every run, and a stale input becomes a diagnosed incident before the first statement goes out.

What stays the same. Kestra, your flows, triggers and namespaces, your dbt project, your warehouse and your on-call rota. For the numbers, see the ROI of agentic data operations.

The sentence for upstairs: "Kestra runs our workflows; Data Workers checks what they produce and gets it fixed with our approval when one goes wrong, so customers see the right number."

Getting started

Start with a pilot. Pick the Kestra flows that feed the numbers leaders and customers read, give Data Workers a Kestra API token and read access to the tables those flows write, and run at L1 for a few weeks: every incident arrives with a cause and a blast radius. Then turn on L2 for one domain, and open L3 for a narrow class once the receipts show the agents were right. The pilot path is on the pricing page, and the pilot is credited in full against the first year.

FAQ

How does Data Workers connect to Kestra? Over Kestra's API or MCP server today, with an API token on Enterprise or Cloud or basic auth on open source. It reads executions, task states and logs. Kestra 2.0 and the patched 1.3 line both expose that REST API.

Kestra 2.0 has AI Copilot, AI Agent tasks and an MCP server. What does Data Workers add? Those help your team author flows, let models call tasks and flows, and let assistants start flows as tools. Data Workers covers the output: whether the data a flow wrote is right, what it reached downstream, how to fix the cause, and the proof afterwards.

We use Cases in Kestra Enterprise. Do they overlap? They work together. Cases track an incident next to the executions behind it inside Kestra. Data Workers diagnoses the data side across the warehouse, dbt, BI and the source, and its receipt can be linked from the case or the alert your on-call team already uses.

Will Data Workers change our flows, triggers or KV pairs? No. In Kestra it reads executions. Flow changes arrive as a YAML diff for the owner to save or merge, and the owner starts reruns and disables or re-enables triggers.

Should we expose our Kestra flows to agents over MCP? Expose the flows you would let a person on your team start without review, on a private server, and keep anything that writes to customers or money behind an approval. Data Workers adds the context and the approval flow around those calls.

Where does our data go? The agents run in your infrastructure with your warehouse credentials and model key. Your data stays in your systems; the hosted Conductor sees workflow metadata only.

Sources

  • •Kestra, homepage: "One orchestrator. Every workflow."; "2,100+ plugins", https://kestra.io/ (checked Oct 3, 2026)
  • •Kestra, llms.txt (core model; "Kestra is API-first: every action available in the UI is also available via REST API"), https://kestra.io/llms.txt (checked Oct 3, 2026)
  • •Kestra releases: v2.0.0 (Sept 7, 2026), v2.0.4 and v1.3.41 (Sept 29, 2026), v1.0.60 (Sept 8, 2026), https://github.com/kestra-io/kestra/releases (checked Oct 3, 2026)
  • •Kestra blog, "In Kestra 2.0, any flow is a tool your agents can call" (Sept 4, 2026), https://kestra.io/blogs/2026-09-04-kestra-2-0-flows-as-agent-tools (checked Oct 3, 2026)
  • •Kestra blog, Kestra 2.0 engine rebuild and workers anywhere (Sept 1, 2026), https://kestra.io/blogs/2026-09-01-kestra20-rebuild-engine (checked Oct 3, 2026)
  • •Kestra docs, Triggers (Schedule, Flow, Webhook, Polling, Realtime, MCP Tool; enabling and disabling triggers), https://kestra.io/docs/workflow-components/triggers, and Realtime trigger (Kafka, Google Pub/Sub), https://kestra.io/docs/workflow-components/triggers/realtime-trigger (checked Oct 3, 2026)
  • •Kestra docs, KV Store, https://kestra.io/docs/concepts/kv-store, Flow revisions, https://kestra.io/docs/concepts/revision, Replay, https://kestra.io/docs/concepts/replay, and the Assert task, https://kestra.io/plugins/core/execution/io.kestra.plugin.core.execution.assert (checked Oct 3, 2026)
  • •Kestra docs, AI Copilot (Ask, Edit, Plan; confirmation), https://kestra.io/docs/ai-tools/ai-copilot (checked Oct 3, 2026)
  • •Kestra docs, AI Agents, https://kestra.io/docs/ai-tools/ai-agents (checked Oct 3, 2026)
  • •Kestra docs, MCP Server and MCP Tool Trigger, https://kestra.io/docs/ai-tools/mcp-server, https://kestra.io/docs/workflow-components/triggers/mcp-tool-trigger (checked Oct 3, 2026)
  • •Kestra docs, Cases (Enterprise), https://kestra.io/docs/enterprise/governance/cases (checked Oct 3, 2026)
  • •Data Workers open-source repository (tool registrations in dw-incidents, dw-quality, dw-context-catalog), 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)