Schedule icon
Webhook icon
Queries icon
Query icon
ChatCompletion icon
OpenAI icon
OutputValues icon
If icon
Write icon
SlackIncomingWebhook icon
Fail icon
Log icon

AI Data Quality Anomaly Advisor and SQL Remediation Generator

Audit tables with DuckDB, diagnose data quality anomalies with AI, and deliver copy-paste cleanup SQL to Slack in Kestra.

Categories
AICoreData

Diagram unavailable

We could not build the topology for this blueprint. The flow itself is valid, use the YAML on the left to run it.

When data quality checks fail in daily ELT/ETL pipelines, debugging is slow and manual: upstream webhook contract changes or ingestion bugs slip in negative order revenues, null customer IDs, or duplicate transaction records. Standard pipeline assertions simply fail the build, forcing on-call data engineers to spend an hour querying staging tables by hand, isolating tainted rows, and hand-writing cleanup SQL scripts.

This blueprint audits incoming staging tables using in-memory DuckDB, extracts the exact records violating quality constraints, and feeds them into Kestra's AI plugin using structured JSON Schema. The AI analyzes the business impact, diagnoses the root cause, generates copy-paste SQL remediation statements (DELETE / UPDATE), saves a downloadable post-mortem report artifact in Kestra storage, and sends an actionable triage card directly to Slack.

How it works

  1. stage_orders_batch (io.kestra.plugin.jdbc.duckdb.Queries) loads the incoming batch of records into an in-memory DuckDB table with zero external database dependencies.
  2. audit_quality_anomalies (io.kestra.plugin.jdbc.duckdb.Query) executes multi-constraint SQL auditing using window functions to detect duplicate primary keys, null references, negative numerical values, and out-of-range future timestamps.
  3. diagnose_anomalies_with_ai (io.kestra.plugin.ai.completion.ChatCompletion) evaluates the anomalous records with strict JSON Schema output, diagnosing the upstream failure mode, assessing business impact, and synthesizing exact copy-paste SQL remediation commands.
  4. consolidate_quality_metrics (io.kestra.plugin.core.output.OutputValues) exposes audit statistics and remediation commands into execution outputs.
  5. evaluate_anomaly_gate (io.kestra.plugin.core.flow.If) branches deterministically:
    • Anomalies Detected: Generates an executive Markdown report artifact (data-quality-triage-report.md) in Kestra storage, alerts Slack with the diagnosis and cleanup SQL, and fails the execution to prevent bad data from reaching the production warehouse.
    • Clean Data: Persists the certified records artifact and logs an all-clear verification.
  6. alert_audit_error (errors block) notifies Slack if the audit pipeline errors.

What you get

  • In-memory data quality auditing with zero infrastructure setup.
  • Automated AI root-cause diagnosis for corrupted data batches.
  • Copy-paste SQL remediation commands delivered straight to Slack.
  • Publication-ready triage Markdown report artifacts stored in Kestra internal storage.

Who it's for

  • Data Engineers and Analytics Engineers managing dbt, DuckDB, Snowflake, or PostgreSQL pipelines.
  • Data Platform teams seeking to reduce MTTR on failed morning ELT loads.
  • FinTech and E-Commerce teams requiring strict data integrity before warehouse merge.

Why orchestrate this with Kestra

Data quality remediation requires combining analytical SQL audits, frontier LLM intelligence, storage persistence, and incident alerting. Kestra coordinates DuckDB, AI models, and Slack into a single auditable, declarative pipeline with complete execution lineage.

Prerequisites

  • OpenAI API key or any OpenAI-compatible LLM endpoint (Groq, Ollama, vLLM).
  • Slack incoming webhook URL for data quality alerts.

Secrets

  • OPENAI_API_KEY: API key for the AI diagnosis task.
  • SLACK_WEBHOOK_URL: Slack incoming webhook endpoint for data quality alerts.
  • DATA_QUALITY_WEBHOOK_KEY: Authentication key for event-driven ELT pipeline triggers.

Quick start

  1. Configure OPENAI_API_KEY and SLACK_WEBHOOK_URL in your Kestra namespace.
  2. Import this blueprint into your Kestra instance.
  3. Click Execute with default sample records to see the automated DuckDB audit and generated remediation SQL in action.
  4. Enable the schedule or trigger the webhook from your upstream ingestion pipeline.

How to extend

  • Point url: "jdbc:duckdb:/path/to/warehouse.db" at a persistent DuckDB database file or MotherDuck.
  • Swap the JDBC connection for PostgreSQL, Snowflake, or BigQuery.
  • Add an automated execution task that directly executes the generated remediation_sql upon manual human confirmation.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.