Product
Product11 min readBy The Data Workers Team

Data Workers + OpenLineage: An OpenLineage Integration That Closes the Loop on Failed Runs

An OpenLineage integration for AI agents: keep your emitters and Marquez, let Data Workers trace failed runs, route fixes for approval and report its own runs as OpenLineage events.

OpenLineage is the open standard that records what ran and what it touched. Data Workers does the part after the red job: it works out why a run failed, proposes the fix to the system that owns it, proves the number is right again with a receipt, and reports its own runs back to your OpenLineage backend. That is the whole shape of this OpenLineage integration.

Your estate probably looks like this. The Airflow OpenLineage provider sends an event for every task. Spark jobs on Databricks run the OpenLineage listener from an init script. dbt runs through dbt-ol, so every model emits a START and a COMPLETE or FAIL. A Flink job or two reports through the Flink connector. All of it lands in Marquez, DataHub or your catalog, and the graph looks great in a demo. Then a FAIL event arrives at 02:10, and someone still has to open three tools to find the column that changed two runs upstream.

Data Workers is the agentic data platform that runs the whole data lifecycle next to the lineage you already collect. Here is what it works from, what it sends back, and what stays where it is.

Key takeaways

  • •OpenLineage stays your lineage standard. Your emitters, your Marquez or DataHub backend and your facets stay as they are, and your team keeps browsing lineage there.
  • •Data Workers builds its own context graph. It reads your dbt manifest, your orchestrator's runs and your warehouse checks, and adds the cross-system hops your team records as notes.
  • •A failed run becomes an incident path. A failure turns into a cause, a blast radius and an owner, across Airflow, Spark, dbt and the warehouse.
  • •Fixes go back through the tool that owns them. Code changes are diffs the owner merges; reruns go through the orchestrator's API after a named approval.
  • •Data Workers reports its own runs. Its verification runs go out as OpenLineage run events to the backend you choose, so the fix shows up in the lineage your team already browses.
  • •Start read-only. Connect dbt and your orchestrator, and see each recent failure with its cause and the fix Data Workers would have proposed.

What OpenLineage does, and why teams keep it

OpenLineage defines a small model, a job, a run and a dataset, and lets every tool extend it with facets: schema, column lineage, data quality assertions, parent run and more. Its own docs keep the scope tight: it "defines the metadata for running jobs and their corresponding events." That focus is why it spread. Schedulers, engines and catalogs can agree on one event format without agreeing on anything else.

Teams keep it because it's neutral, open (Apache 2.0, an LF AI & Data Foundation Graduate project) and already wired into the tools they run. Here is where the ecosystem stands this month.

AreaWhat shipsStatus, October 2026
Core specJob, run and dataset events; facets; START, RUNNING, COMPLETE, FAIL, ABORTSpec 2-0-2
Release lineClients and integrations1.53.0, Sept 1, 2026
New in 1.53.0Explicit lineage facets, Spark 4.2 and Python 3.14, Spark column-lineage transformation descriptionsSept 1, 2026
Column lineage facetDIRECT (identity, transformation, aggregation) and INDIRECT (join, filter, group by, sort, window, condition)Facet 1-2-0
Airflowapache-airflow-providers-openlineage, Airflow 2.11 and later2.20.2, Sept 29, 2026
SparkOpenLineageSparkListener; on Databricks it installs through a cluster init scriptAvailable
dbtdbt-ol, a drop-in wrapper for the dbt commandAvailable
FlinkOne connector: 1.x needs a job-code change; 2.x uses Flink's native lineage interfaces and supports Flink SQLAvailable
DagsterCommunity-maintained dagster-openlineage package, listed in Dagster's docsv0.2.1, May 2026
MarquezThe reference implementation of the OpenLineage API, an LF AI & Data projectAvailable

The receiving side has grown fast too. Snowflake external lineage accepts OpenLineage COMPLETE events at a REST endpoint (GA since September 3, 2026, Enterprise Edition). Google Cloud's Knowledge Catalog (formerly Dataplex Universal Catalog) consumes events through ProcessOpenLineageRunEvent. Microsoft Purview takes Airflow lineage through OpenLineage and Azure Event Hubs, in preview. DataHub has an OpenLineage REST endpoint and a Spark listener plugin. Atlan has OpenLineage connectors for Airflow, MWAA, Astronomer, Cloud Composer, Spark, Flink and any generic emitter. Databricks records its own lineage in Unity Catalog; Spark jobs there report OpenLineage through the listener.

What Data Workers works from, and what it sends back to OpenLineage

OpenLineage keeps the record of what ran; your backend keeps it browsable. Data Workers works next to it from a governed context graph in Data Context Wizard, built from your dbt project, your orchestrator's runs and your warehouse checks. Hops the dbt manifest can't see, like a Spark task that reads a Postgres table over JDBC, go into the graph as notes your team records once, with an author and a source. When Data Workers runs a verification, it reports that run as an OpenLineage run event, with its inputs and outputs, to the endpoint you configure.

Data Workers works fromData Workers sends back
The dbt manifest: models, sources, tests and their dependenciesA blast radius before a change ships, with owners
Orchestrator runs: DAG runs and task instances, with their stateDiffs for the repo that owns the fix
Warehouse checks: quality checks and metric baselines on the tables that matterReruns through Airflow, Dagster or Prefect, after a named approval
Team notes in the context graph: cross-system hops, readers and ownersIts own run events, sent to the OpenLineage backend you choose
The pull requests behind your jobs and modelsA receipt for every change, in Spellbook
What Data Workers reads from OpenLineage and what it writes back through OpenLineage

A few tools do most of the work. trace_cross_platform_lineage walks lineage upstream and downstream across the context graph. blast_radius_analysis lists the models, metrics and recorded readers a change reaches, with owners and severity. The Autonomous Data-Conductor sequences the run and holds every change at the autonomy level you set. Data Workers reports its verification runs to your OpenLineage endpoint, so Marquez shows the verification run next to the jobs it checked.

One incident, through OpenLineage

Here is a night most OpenLineage teams will recognize. It's an illustration, not a customer case.

At 00:40 an application release changes orders.discount in Postgres from numeric to text, so a handful of rows now say N/A. At 01:05 the Airflow DAG orders_daily runs a Spark task on Databricks that reads orders over JDBC and writes lake.orders_clean. The Spark listener emits a COMPLETE event, and its schema facet now shows discount as a string. Nothing fails. At 02:10 the next task runs dbt-ol build, and fct_revenue on Snowflake emits a FAIL event: numeric value 'N/A' is not recognized. Marquez shows a red job. Finance opens the revenue dashboard at 08:30.

StepWhere it runsWhat happensWho decides
1. DetectAirflowData Workers reads the orders_daily run and sees the dbt-ol build task fail on fct_revenue: numeric value 'N/A' is not recognized. An incident opens.Data Workers, read-only
2. Diagnosedbt, Postgresdiagnose_incident ranks the cause: the failing cast is on discount, and the 00:40 release recorded against orders is the only recent upstream change.Data Workers, read-only
3. TracePostgres, Spark, dbttrace_cross_platform_lineage follows fct_revenue through the dbt sources to orders_clean, and through the team's note on the Spark task back to orders in Postgres. The red job in Marquez shows the same path.Data Workers, read-only
4. Blast radiusdbt, Snowflake, BIblast_radius_analysis lists two dbt models, the net_revenue metric and, from the team's note, one revenue dashboard, with their owners.Data Workers, read-only
5. ProposeGitHubData Workers proposes a diff on the Spark job's repo: parse discount explicitly, send non-numeric values to a quarantine column, and add a schema expectation. The diff carries the path, the blast radius and the rollback. It also notifies the Postgres table's owner.Data Workers proposes
6. ApproveGitHub or SpellbookThe on-call data engineer reviews and merges at 07:10. CI runs as usual.A named engineer
7. RerunAirflowData Workers queues a new run of orders_daily through the Airflow 2 REST API, the one DAG the approval covers. On Airflow 3 the owner triggers the approved run.Approved in step 6
8. VerifySnowflake, MarquezBoth tasks complete, and their COMPLETE events land in Marquez as usual. Data Workers reads the run, re-runs its null and row-count checks on fct_revenue, confirms net revenue is back on its monitor_metrics baseline, records the receipt, and reports its verification run to Marquez as an OpenLineage run event.Data Workers, read-only
Incident timeline across the stack: what OpenLineage, your team and Data Workers each do, step by step

Every handoff went through the tool that owns it. Data Workers supplied what no single event carries: the ranked cause, the path across five systems, the proposed diff with its blast radius, the approval, and the checks that show revenue back on its baseline. For the Airflow side of this wiring in detail, read Data Workers + Airflow; for Dagster and Prefect, read Data Workers + Dagster and Prefect.

Why doesn't OpenLineage just do this itself?

OpenLineage built exactly the right thing for its job: a neutral record of what ran and what it touched. Its scope is the metadata for running jobs, sent to a backend you choose. That narrow scope is the reason Airflow, Spark, dbt, Flink, Snowflake and Google Cloud all adopted it. A standard that every emitter and every consumer trusts has to stay out of anyone's production systems.

Acting on lineage is a different product category. It needs approvals from the owners of systems the standard only describes, writes to a Spark repo and a dbt project in the same incident, rollback for each step, receipts an auditor can read, context from sources that don't emit events, and accountability for changes in tools the OpenLineage project doesn't run. Putting that into the spec would turn a shared standard into a vendor product that every integration had to trust. The acting layer on top is the product Data Workers is.

How to connect today

Data Workers is MCP-native: each agent is an MCP server, so your coding agent (Claude Code, Codex or Cursor) can ask about a failed run, hand off a fix or read a receipt from the same session. Your emitters don't change.

Week one: read only.

  • •Connect the dbt project and your orchestrator read-only: the dbt manifest, and the DAG runs and task instances from Airflow, Dagster or Prefect.
  • •Record the hops dbt can't see as context-graph notes: which Spark task reads which source, which dashboard reads which mart. Your Marquez graph is the quickest place to list them.
  • •Give Data Workers read access to the Git repos behind your jobs and models.
  • •Every agent starts observe-only. The first thing you see is each recent failed run with its cause, its blast radius and the fix it would have proposed.

Example: the trace and blast radius calls

[
  {
    "tool": "blast_radius_analysis",
    "arguments": {
      "assetId": "orders_clean"
    }
  },
  {
    "tool": "trace_cross_platform_lineage",
    "arguments": {
      "assetId": "fct_revenue",
      "direction": "upstream",
      "maxDepth": 6,
      "includeOrchestration": true
    }
  }
]

Week two onward: proposals, then approved runs.

  • •Turn on blast-radius comments on pull requests that change a dataset other jobs read.
  • •Let Data Workers propose diffs for the owning repo, with no permission to merge.
  • •Allow reruns of specific DAGs, jobs or deployments, each behind an approval.
  • •Point the context agent at your OpenLineage endpoint if you want Data Workers' own runs reported to your backend.

The scopes above are an illustration; use the permission sets your tools offer. If you run Snowflake and Databricks side by side, the cross-cloud lineage how-to shows how the same graph spans both platforms.

How it fits together

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

Your team works in its coding agent and reviews in Spellbook Data Catalog, which is in preview. The Data-Agents Swarm does the work, and the Conductor runs each fix end to end. OpenLineage stays your collection standard, and Marquez, DataHub or your catalog stays where people browse lineage. Nothing migrates. Data Workers stores metadata and scrubbed facts, not copies of your tables.

Guardrails: approvals, and what OpenLineage owns

What stays with OpenLineage and your backend. The event format, the facets, your emitters and their configuration, the backend that stores events, and the lineage UI your team already uses. Data Workers adds its own run events next to yours; it never edits stored events.

What Data Workers enforces on top.

  • •Read-only start. New deployments are observe-only. You extend autonomy one domain at a time as the receipts earn trust.
  • •Autonomy per domain, L0 to L4. Blast-radius comments can run freely while Spark job changes stay at "propose".
  • •Approvals where they belong. Code changes are diffs your reviewers merge, so branch protection and your reviewers decide. Reruns name the one DAG or job an approval covers.
  • •No self-approval. No agent can promote its own work.
  • •Receipts. Every change records the failure that started it, the path, the diff, the approver, the checks run, the before and after values, and the rollback path in a tamper-evident, hash-chained log.
  • •Rollback. A merged fix reverts like any commit, and the receipt records how.
  • •Least privilege. Data Workers acts with the grants you give it, through each tool's own permission system.
The autonomy ladder: L0 manual, L1 observe, L2 propose, L3 act reversibly, L4 autonomous

What changes for your team

Data engineers stop clicking from a red job in Marquez back through three runs. They open a proposed diff that already names the column, the change that broke it and every recorded reader downstream. The on-call rotation changes from "find it" to "approve it". Platform leads see every fix land in the lineage backend they spent a year instrumenting, next to the jobs it checked. Governance gets a record of every change tied to the runs around it. For a deeper look at impact analysis on lineage, see lineage agent impact analysis and column-level lineage; for adding OpenLineage to pipelines you're writing with a coding agent, see OpenLineage instrumentation with Claude Code.

The case for your CFO

The outcome. A failed run stops being a red box someone decodes by hand. It becomes a reviewed, verified fix before the business opens the dashboard, with a record of what changed and who approved it.

The risk story. At the observe level, agents read runs and metadata and report; they change nothing. At the propose level, they propose diffs and a named engineer merges. Acting levels are reversible and only for domains you choose. No agent can promote its own work, and every change has a rollback path. Each receipt holds the trigger, the diff, the approver, the blast radius, the checks run and the before and after values. There is no migration: your emitters, backend and catalog stay as they are.

Why now. In 2026 the warehouses and clouds started ingesting OpenLineage natively, and 1.53.0 added explicit lineage facets. Collection is solved. The bottleneck is acting on the failures it records.

The first win. Connect dbt and your orchestrator read-only and get each recent failed run with its upstream cause, its blast radius and the fix Data Workers would have proposed.

What stays the same. OpenLineage, Marquez or DataHub, Airflow, Spark, dbt, your catalog and the coding agent your engineers already use.

The pilot path. Start with a pilot on one pipeline domain, read-only first, then proposed diffs, then approved reruns. The pilot is credited in full against the first year.

One sentence for upstairs: "Every pipeline already reports what it ran in OpenLineage; Data Workers turns its failures into fixes our owners approve, with a receipt for each one."

When OpenLineage on its own is enough

If one team owns every pipeline, failures are rare and the person on call can read the Marquez graph and fix things by hand, OpenLineage and a good backend carry you. Once a failure regularly crosses Airflow, Spark, dbt and the warehouse, or more than one team owns the path, an operating layer that traces the failure, proposes the fix and verifies it starts paying for itself.

FAQ

Where does Data Workers get lineage from? From its own context graph: the dbt manifest, your orchestrator's runs, your warehouse checks and the cross-system hops your team records as notes. Your OpenLineage backend stays the place people browse the full event history.

Do we have to replace Marquez or DataHub? No. Keep your backend and its UI. Data Workers adds the investigation, the approvals and the receipts next to it, and reports its own runs into it.

Does Data Workers emit OpenLineage events? Yes. Data Workers reports its own runs, with inputs and outputs, as OpenLineage run events to the endpoint you configure, so its verification runs appear in your backend next to the jobs they checked.

What about hops dbt can't see? A Spark task reading Postgres over JDBC, or a dashboard reading a mart, goes into the context graph as a note your team records once, with an author and a source. The trace and the blast radius use it from then on.

Does Data Workers need write access to our pipelines? Not to start. It reads runs and metadata. Fixes arrive as diffs your reviewers merge, and reruns go through the orchestrator's API only after a named approval.

Sources

Sources for OpenLineage capabilities and statuses, checked October 2, 2026: About OpenLineage (model, scope, integrations, Marquez, LF AI & Data Graduate, docs version 1.53.0), OpenLineage releases (1.53.0, Sept 1, 2026) and changelog (explicit lineage facets; Dagster integration moved out of the repo), the core spec (2-0-2), the column lineage facet, the Airflow OpenLineage provider (2.20.2), the Spark integration and OpenLineage on Databricks (init script), the dbt integration (dbt-ol), the Flink integration, Dagster and OpenLineage, Marquez, the OpenLineage ecosystem (consumers and producers), Snowflake external lineage, Google Cloud Knowledge Catalog OpenLineage integration, Microsoft Purview lineage from Airflow (preview), DataHub OpenLineage and Atlan's connector index. Product names and statuses change quickly; if we've got something wrong, tell us and we'll fix it.