A failed pipeline rarely provides all the information needed to fix it. The error might appear in Airflow, the failing SQL might come from dbt, and the change that caused it might be in an ingestion connector or source application.
AI agents for data engineering combine language models with access to code, metadata, execution history, and controlled tools. Their useful work includes developing transformations, investigating failures, reviewing schema changes, planning backfills, and checking whether a repair actually restored the intended data.
The difficult part is connecting those activities. Writing a corrected query does not establish which partitions need rebuilding. A successful backfill does not establish that downstream metrics are correct. And a convincing diagnosis does not authorize a production change.
This guide examines those boundaries across dbt, Airflow, Dagster, ingestion systems, and cloud warehouses, including how Claude Code and Codex fit into the workflow.
What can AI agents do across the data engineering stack?
Start with a task that has an observable result. “Improve data reliability” is too broad to evaluate. “Investigate this failed execution and produce a supported diagnosis” gives the agent a defined scope and gives the engineer something to verify.
Develop and debug dbt models
Model development is a useful starting point: generating transformations, adding tests, documenting dependencies, and reviewing changes to existing SQL. The dbt MCP server exposes development, metadata, and execution capabilities that can support these tasks, depending on configuration and platform access.
The request still needs a data contract. For an orders model, specify whether the output contains one row per order, one row per order line, or a history of order changes. Include how refunds, deletions, and late updates should behave.
Consider this deliberately flawed filter inside an incremental model:
{% if is_incremental() %}
where created_at > (
select max(created_at)
from {{ this }}
)
{% endif %}
Assume the target already contains records. An order created last month and refunded today will be excluded when its created_at remains unchanged. Configuring order_id as the unique_key does not help if the updated record never enters the incremental batch. dbt’s filtering logic determines the incoming rows; its incremental strategy determines how those rows are applied.
A useful agent should identify that distinction, trace a missed update from source to target, and reproduce the behavior. Replacing created_at with updated_at is only a candidate repair. The source must reliably update that field, and the design still needs to handle timestamp ties, late delivery, replay, and deletions.
Investigate Airflow failures and plan safe retries
For an Airflow investigation, require the exact DAG run, task, attempt, relevant map index, data interval, and deployed code version. Then distinguish failure during execution from waiting for capacity, missing upstream input, or processing the wrong interval.
The recovery action depends on that distinction. A temporary connection failure may justify retrying. Invalid SQL will not become valid on another attempt. A task that partially committed output needs an examination of its write behavior before replay.
Airflow recommends transaction-like tasks whose reruns produce the same outcome, with reads and writes tied to specific partitions rather than whichever data happens to be “latest.”
Ask the agent for a scoped replay plan: the task instances to clear, the intervals to process, the output already written, and the checks required afterward.
For the incremental-model example, rerunning the task unchanged would continue missing the refund. Operational recovery requires fixing the transformation, not simply restoring a green task status.
Diagnose Dagster assets and partition dependencies
In Dagster, organize evidence around asset keys, partition keys, materializations, and asset checks. Asset checks validate properties of assets; appropriately configured blocking checks can prevent dependent execution when a check fails.
An agent can investigate why a partition is absent, why a check failed, or whether an upstream correction requires downstream rematerialization. Partition mappings are important here: a downstream partition can depend on multiple upstream partitions rather than just the same date.
Suppose an output for each day sums that day and the previous six days. Correcting an input on September 1 can affect outputs from September 1 through September 7. Rebuilding only the September 1 output leaves the rest potentially incorrect.
Require the agent to derive recovery scope from transformation semantics and partition dependencies. A list of downstream asset names is not yet a backfill plan.
Investigate ingestion, CDC, and streaming state
For Airbyte, Fivetran, or a custom ingestion service, the first question should be whether the relevant change reached the destination at all. Request connector status, schema-change history, cursor or offset information, and the affected delivery interval.
Schema problems do not always produce downstream SQL errors. Airbyte, for example, pauses connections when an existing cursor or primary key is removed. The downstream symptom may therefore be stale data while the destination schema remains queryable.
For Debezium, Kafka, Spark, or Flink workflows, make replay semantics part of the investigation. Establish what the source retains, which offsets or checkpoints identify progress, and how the sink recognizes records already written.
Spark’s foreachBatch provides at-least-once writes by default; its batch identifier can support deduplication when the application implements it. Restarting a streaming job therefore does not, by itself, establish duplicate-free recovery.
An agent should explain these consequences before recommending a reset, resnapshot, or checkpoint change.
Explain quality failures and metric discrepancies
Keep deterministic checks in tools built to execute them. Great Expectations Checkpoints, for example, run validations and actions based on their results. An agent can investigate those results rather than replacing the assertions with a model’s opinion.
For a revenue discrepancy, compare the underlying populations before comparing totals. Are cancelled orders included? Does a payment-attempt join multiply orders? Do the reports use the same currency conversion and day boundary?
Retrieve approved metric definitions where they exist. The dbt Semantic Layer centralizes metric definitions and associated query logic, providing a better starting point than asking an agent to infer “revenue” from column names.
The useful output is a specific explanation and a regression check not a rewritten metric that happens to agree with one dashboard.
Warehouse-specific problems need warehouse-specific evidence
A slow pipeline does not necessarily need a SQL rewrite. Separate time spent waiting, time spent executing, excessive resource consumption, and incorrect results before proposing a change.
Snowflake: queueing, spilling, and query behavior
Snowflake query history exposes execution time, queueing, transaction blocking, scan metrics, and spilling. These distinguish a query waiting behind other workloads from one doing expensive work after it starts. The Account Usage view can lag by up to 45 minutes, so an immediate investigation also needs current execution evidence.
For queueing, inspect overlapping workloads and warehouse configuration. For unexpectedly large intermediate results, inspect join cardinality and filtering. Compare equivalent executions before attributing a regression to a deployment.
Be particularly careful with incremental optimizations. dbt’s incremental_predicates can become part of the target-side merge condition. Excluding older target rows may prevent an incoming update from matching an existing record. A scan reduction can therefore change correctness, not just performance.
The agent should demonstrate that required records remain reachable before benchmarking the optimization.
BigQuery: scan costs and investigation budgets
Give diagnostic queries explicit cost boundaries. BigQuery supports dry-run estimates and maximum-bytes-billed controls for on-demand queries. Adding LIMIT 100 is not a general substitute: on non-clustered tables, limiting returned rows does not reduce the data read.
When reviewing an expensive transformation, ask which columns and partitions are necessary for the business interval. A proposed date restriction must still capture required historical corrections.
Evaluate the change against equivalent inputs. A faster query over fewer required records is an incomplete computation, not a successful optimization.
Databricks and Spark: skew versus memory pressure
Inspect stage and task evidence rather than total runtime alone. Databricks documents skew as uneven work across tasks and spilling as execution-memory pressure that moves data to disk. These require different investigations.
A few unusually long tasks justify examining hot keys and partition distribution. Broad spilling justifies examining intermediate data size, joins, and available memory. “Add more workers” should not be the default response to both.
Require a proposed change to preserve results and to explain which observed bottleneck it addresses.
Redshift: distinguish waiting from execution
Redshift’s SYS_QUERY_HISTORY separates queue time from execution time and includes statuses and errors. Use that evidence before recommending workload-management changes or query rewrites.
Across warehouses, preserve the comparison conditions: input interval, volume, compute configuration, concurrency, and cache behavior. Otherwise an agent may take credit for an improvement caused by an easier workload.
How AI agents should investigate root causes
A practical investigation needs three things: a specific execution, a way to retrieve related evidence, and a stopping condition.
Establish which execution is being examined
“This model failed” is insufficient. The same model can run under different commits, roles, source states, and data intervals.
Preserve an execution record that connects those identities. For dbt, run_results.json refers to resources through unique_id, linking execution results to definitions in manifest.json. Archive artifacts from the relevant invocation rather than relying on files produced by a later run.
An illustrative investigation record might contain:
environment: production
asset: model.analytics.fct_orders
orchestrator_run:
task_attempt: 2
data_interval:
deployed_commit:
dbt_invocation:
warehouse_query:
evidence_observed_at:
Retrieve connected evidence
Start with the failed node or violated invariant. Retrieve its compiled SQL, immediate dependencies, relevant logs, schema evidence, and last comparable successful execution. Expand only when a finding justifies it.
Use code for exact operations: resolving identifiers, traversing known dependencies, comparing schemas, and grouping duplicate notifications. Use the language model to interpret findings and choose useful follow-up checks.
Preserve provenance. A current schema and yesterday’s failed query may describe different states. When a historical snapshot is unavailable, say so instead of silently substituting today’s schema.
Test explanations against the evidence
Consider a hypothetical model that fails after a column rename. The rename is a candidate cause.
To support it, establish that the renamed source was visible to the failing execution, that the executed SQL referenced the old identifier, and that the timing fits. Check alternatives such as the wrong database target or a permission change.
A recent deployment is not sufficient evidence by itself. Neither is finding an error message that resembles a past incident.
The report should distinguish observed facts, the best-supported explanation, and unverified assumptions. Where possible, reproduce the failure in isolation and show that the proposed change removes it without breaking independent checks.
How Claude Code and Codex fit into data engineering workflows
Claude Code provides a tool-using loop for gathering context, making changes, and verifying results. Codex supports non-interactive execution through codex exec, including use in scripts and CI. Claude Code also provides programmatic execution and structured output.
These capabilities can supply the investigation and coding component of a larger workflow. They do not require the engineer to paste every log into a conversation.
A service can receive an incident, resolve its execution identifiers, prepare the relevant checkout, and expose diagnostic tools. The agent then inspects evidence, requests additional checks, and proposes a change.
Repository instructions supply conventions
Use CLAUDE.md for Claude Code and AGENTS.md for Codex to describe stable project rules: model grain, naming conventions, approved development environments, required tests, and actions that need approval.
These files should explain how the project is supposed to work. They should not masquerade as an authoritative record of current warehouse state.
For example, document that orders can change after their creation date. Retrieve the actual missed orders and deployed incremental predicate during the investigation.
MCP supplies access, not operational understanding
Both tools support MCP for connecting external capabilities. The integration can expose warehouse diagnostics, orchestrator records, dbt metadata, and other approved systems.
Prefer narrow operations over one unrestricted production connection. An illustrative interface could be:
get_execution_evidence(execution_ref)
get_schema_evidence(relation_ref, as_of)
run_diagnostic(template_id, parameters, budget)
validate_patch(patch_ref, fixture_ref, environment)
The service behind each operation should enforce scope and permissions. An operation asking for a historical schema must return “unavailable” when no historical evidence exists.
The dbt MCP server already exposes tools for node details, lineage, run artifacts, and development commands. Review which are enabled: metadata lookup and model execution have different consequences
Challenges of using Claude Code and Codex for data engineering
The main limitations appear where repository work meets changing production state.
Incomplete context, and context from the wrong incident
An agent may have the SQL but not the deployed version, the model lineage but not the affected partition, or the current schema but not the schema used by the failed run.
Adding more files does not necessarily fix this. The missing information may be the relationship between two records. Maintain execution mappings and timestamped evidence outside the conversation. Also distinguish “no matching record,” “permission denied,” and “evidence expired.” Those outcomes should lead to different conclusions.
Long conversations introduce another constraint: Claude Code documents that context compaction can discard earlier details. Durable incident state should therefore live in the surrounding application, not depend on the agent remembering everything.
No representative test environment
The problem is not that coding agents cannot execute tests. It is that the available environment may not reproduce the behavior under investigation.
For an incremental repair, an empty target exercises insertion but not updates to existing state. dbt explicitly notes that incremental unit tests validate the rows to be inserted or merged, not the final target table after the operation.
Provide a stateful integration fixture with a new record, an update to an older record, late delivery, repeated input, and out-of-order versions where relevant. Run the actual materialization and inspect the resulting target. Match the warehouse dialect, adapter version, macros, and relevant permissions. Use separate representative workloads for performance testing. Without that environment, the agent can produce a proposal, but the unperformed checks must remain visible. Production should not become the test environment by default.
Proactivity gap & unattended execution
It would be inaccurate to describe these tools as incapable of automation, it's just can not make proactive agent in one shot. Codex can run within scheduled scripts or CI. Claude Code routines support schedules, API calls, and GitHub events; routines are currently documented as a research-preview feature.
What still needs engineering is the data-specific monitoring loop: which signal matters, which execution it concerns, which related alerts belong together, and what proves recovery.
A recurring repository review is not equivalent to monitoring whether ingestion will miss a downstream freshness deadline. The latter needs operational state and a defined service requirement.
Tool access can exceed the intended authority
A local coding sandbox and a remote warehouse role protect different boundaries. Apply both. OpenAI documents sandbox, approval, network, and audit controls for Codex; dbt separately warns that MCP-exposed commands can modify warehouse objects.
Use separate credentials for investigation, development validation, and production execution. Restrict returned data as well as writes: read-only access can still disclose sensitive records.
Require consequential actions to name the environment, affected objects, partitions, preconditions, and approval. Treat retrieved logs and documents as evidence, not instructions that can expand those permissions.
How to evaluate AI agents for data engineering
Begin with historical incidents and provide only evidence available at the simulated investigation time. Exclude postmortems and later patches that reveal the answer.
Compare the agent-assisted workflow with the existing runbook. Measure supported diagnoses, incorrect causal claims, successful repairs, reviewer effort, and combined model-and-warehouse cost. Include incidents with insufficient evidence, where requesting a missing record is better than inventing a cause.
Start with evidence collection and read-only investigation. Permit development patches once diagnoses are dependable. Grant narrowly defined production actions only when their preconditions, approval, and verification are established.
For the missed-refund example, success means finding why the update was excluded, demonstrating the behavior, repairing historical records, and confirming that replay preserves the intended state.
Measure your agent's reliability by using publicly available data engineering benchmarks for agents.
Where Upriver fits into this architecture
The workflows above expose one of the harder problems with AI agents for data engineering: the model is rarely the main bottleneck. The bottleneck is assembling enough reliable context for the agent to understand what happened across the data stack.
That is the problem Upriver is designed around.
Upriver builds a continuously updated context layer across the systems involved in data engineering warehouses, dbt projects, orchestrators, lineage, code, metadata, and documentation and uses that context to support agent-driven DataOps workflows.
Instead of starting an investigation from a single Airflow error or dbt failure, the agent can reason across related evidence.
The important part is not simply collecting these objects. The system has to preserve the relationships between them: which dbt invocation belongs to which orchestrator run, which query was generated by that execution, which schema was visible at the time, and which downstream assets depend on the failing model.
That gives an agent a much stronger starting point for root-cause analysis.
That is also where AI agents for data engineering become materially different from running a general-purpose coding agent against a data repository. Claude Code or Codex can provide strong reasoning and code-generation capabilities.
Upriver is focused on that layer: giving data engineering agents enough context to investigate production problems and generate bounded, reviewable actions rather than reasoning from isolated logs or source code.