id: ai-data-quality-anomaly-advisor
namespace: company.data
description: |
Audit inbound staging tables for data quality anomalies using DuckDB,
evaluate corrupt records with Kestra's AI plugin using strict JSON Schema,
and synthesize copy-paste SQL remediation commands for on-call data engineers.
triggers:
- id: scheduled_quality_audit
type: io.kestra.plugin.core.trigger.Schedule
description: Daily pre-ETL data quality audit. Shipped disabled by default.
cron: "0 6 * * 1-5"
disabled: true
- id: webhook_quality_audit
type: io.kestra.plugin.core.trigger.Webhook
description: Webhook triggered after upstream batch or file arrival in data staging.
key: "{{ secret('DATA_QUALITY_WEBHOOK_KEY') }}"
inputs:
- id: table_name
type: STRING
defaults: raw_orders_staging
description: Staging table or dataset name under evaluation.
- id: max_tolerable_anomalies
type: INT
defaults: 0
description: Maximum allowed anomalous rows before halting downstream warehouse loading.
- id: model_name
type: STRING
defaults: gpt-4o-mini
description: LLM used for automated data quality root cause diagnosis and SQL synthesis.
- id: staging_orders
type: JSON
defaults:
- order_id: "ORD-1001"
customer_id: "CUST-452"
order_amount: 149.99
order_date: "2026-10-04 10:15:00"
status: "COMPLETED"
- order_id: "ORD-1002"
customer_id: null
order_amount: 89.50
order_date: "2026-10-04 10:20:00"
status: "COMPLETED"
- order_id: "ORD-1003"
customer_id: "CUST-881"
order_amount: -45.00
order_date: "2026-10-04 10:25:00"
status: "REFUNDED"
- order_id: "ORD-1004"
customer_id: "CUST-104"
order_amount: 520.00
order_date: "2099-01-01 00:00:00"
status: "PENDING"
- order_id: "ORD-1001"
customer_id: "CUST-452"
order_amount: 149.99
order_date: "2026-10-04 10:15:00"
status: "COMPLETED"
description: Candidate staging records evaluated for business rule anomalies.
tasks:
- id: stage_orders_batch
type: io.kestra.plugin.jdbc.duckdb.Queries
description: Loads incoming batch records into an in-memory DuckDB staging table.
url: "jdbc:duckdb:"
fetchType: NONE
sql: |
CREATE TABLE raw_orders_staging (
order_id VARCHAR,
customer_id VARCHAR,
order_amount DOUBLE,
order_date TIMESTAMP,
status VARCHAR
);
{% for row in inputs.staging_orders %}
INSERT INTO raw_orders_staging VALUES (
'{{ row.order_id }}',
{% if row.customer_id is null %}NULL{% else %}'{{ row.customer_id }}'{% endif %},
{{ row.order_amount }},
TIMESTAMP '{{ row.order_date }}',
'{{ row.status }}'
);
{% endfor %}
- id: audit_quality_anomalies
type: io.kestra.plugin.jdbc.duckdb.Query
description: Executes multi-rule audit detecting null references, negative
revenue, duplicate IDs, and future dates.
url: "jdbc:duckdb:"
fetchType: FETCH
sql: |
WITH ranked_orders AS (
SELECT
order_id,
customer_id,
order_amount,
order_date,
status,
ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY order_date) AS row_num
FROM raw_orders_staging
)
SELECT
order_id,
customer_id,
order_amount,
strftime(order_date, '%Y-%m-%d %H:%M:%S') AS order_date,
status,
CASE
WHEN customer_id IS NULL THEN 'NULL_CUSTOMER_REFERENCE'
WHEN order_amount < 0 THEN 'NEGATIVE_ORDER_AMOUNT'
WHEN order_date > TIMESTAMP '2026-12-31 23:59:59' THEN 'FUTURE_ORDER_TIMESTAMP'
WHEN row_num > 1 THEN 'DUPLICATE_ORDER_ID'
ELSE 'VALID'
END AS anomaly_rule
FROM ranked_orders
WHERE customer_id IS NULL
OR order_amount < 0
OR order_date > TIMESTAMP '2026-12-31 23:59:59'
OR row_num > 1;
- id: diagnose_anomalies_with_ai
type: io.kestra.plugin.ai.completion.ChatCompletion
description: Analyzes anomalous data rows, explains business impact, and
synthesizes copy-paste cleanup SQL.
provider:
type: io.kestra.plugin.ai.provider.OpenAI
apiKey: "{{ secret('OPENAI_API_KEY') }}"
modelName: "{{ inputs.model_name }}"
messages:
- type: SYSTEM
content: |
You are a Staff Data Quality Engineer and Database Reliability Architect.
Your task is to analyze anomalous records from a failed data quality scan against table: {{ inputs.table_name }}.
Explain the root cause hypothesis, determine downstream business impact, and generate
an exact, copy-paste SQL cleanup script (e.g. DELETE or UPDATE statements with WHERE clauses)
that an on-call engineer can immediately execute.
Return valid JSON strictly matching the provided schema.
- type: USER
content: |
Table under audit: {{ inputs.table_name }}
Anomalous records detected:
{{ outputs.audit_quality_anomalies.rows | json }}
configuration:
responseFormat:
type: JSON
jsonSchema:
type: object
properties:
anomalies_detected:
type: boolean
detected_rules_breached:
type: array
items:
type: string
business_impact:
type: string
root_cause_diagnosis:
type: string
remediation_sql:
type: string
quarantine_recommended:
type: boolean
required:
- anomalies_detected
- detected_rules_breached
- business_impact
- root_cause_diagnosis
- remediation_sql
- quarantine_recommended
- id: consolidate_quality_metrics
type: io.kestra.plugin.core.output.OutputValues
description: Exposes structured data quality metrics and remediation script into
execution outputs.
values:
anomalous_rows_count: "{{ outputs.audit_quality_anomalies.rows ?
outputs.audit_quality_anomalies.rows | length : 0 }}"
anomalies_detected: "{{ outputs.diagnose_anomalies_with_ai.jsonOutput.anomalies_detected }}"
breached_rules: "{{
outputs.diagnose_anomalies_with_ai.jsonOutput.detected_rules_breached
}}"
quarantine_recommended: "{{
outputs.diagnose_anomalies_with_ai.jsonOutput.quarantine_recommended }}"
remediation_sql: "{{ outputs.diagnose_anomalies_with_ai.jsonOutput.remediation_sql }}"
- id: evaluate_anomaly_gate
type: io.kestra.plugin.core.flow.If
description: Quarantines batch and alerts on-call if anomalous records exceed threshold.
condition: "{{ outputs.consolidate_quality_metrics.anomalous_rows_count >
inputs.max_tolerable_anomalies }}"
then:
- id: export_triage_report
type: io.kestra.plugin.core.storage.Write
description: Persists executive data quality triage report into Kestra internal
storage.
fileName: "data-quality-triage-report-{{ inputs.table_name }}.md"
content: |
# Data Quality Anomaly Triage & Remediation Report
- **Table Audited:** `{{ inputs.table_name }}`
- **Audit Timestamp:** `{{ now() }}`
- **Total Anomalous Rows:** `{{ outputs.consolidate_quality_metrics.anomalous_rows_count }}`
- **Rules Breached:** `{{ outputs.consolidate_quality_metrics.breached_rules | json }}`
- **Quarantine Recommended:** `{{ outputs.consolidate_quality_metrics.quarantine_recommended }}`
## Root Cause Diagnosis
{{ outputs.diagnose_anomalies_with_ai.jsonOutput.root_cause_diagnosis }}
## Downstream Business Impact
{{ outputs.diagnose_anomalies_with_ai.jsonOutput.business_impact }}
## Recommended SQL Cleanup Commands
Execute the following SQL commands to clean or quarantine the affected staging records:
```sql
{{ outputs.consolidate_quality_metrics.remediation_sql }}
```
- id: notify_data_engineers_slack
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Delivers triage card to Slack with root-cause summary and
copy-paste remediation SQL.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "🚨 *DATA QUALITY ALERT: {{ inputs.table_name }}*\n*Anomalous Rows:* `{{ outputs.consolidate_quality_metrics.anomalous_rows_count }}`\n*Breached Rules:* `{{ outputs.consolidate_quality_metrics.breached_rules | json }}`\n*Diagnosis:* {{ outputs.diagnose_anomalies_with_ai.jsonOutput.root_cause_diagnosis }}\n*Copy-Paste Remediation SQL:*\n```sql\n{{ outputs.consolidate_quality_metrics.remediation_sql }}\n```\n*Triage Report:* `{{ outputs.export_triage_report.uri }}`"
}
- id: fail_quarantined_batch
type: io.kestra.plugin.core.execution.Fail
description: Fails execution to halt downstream warehouse ingestion until cleaned.
errorMessage: "Table {{ inputs.table_name }} failed quality audit: {{
outputs.consolidate_quality_metrics.anomalous_rows_count }} anomalous
rows detected."
else:
- id: export_clean_dataset
type: io.kestra.plugin.core.storage.Write
description: Persists verified clean records artifact for warehouse loading.
fileName: "certified-clean-{{ inputs.table_name }}.json"
content: "{{ inputs.staging_orders | json }}"
- id: log_clean_verification
type: io.kestra.plugin.core.log.Log
description: Records clean table verification.
message: "Data quality audit passed for {{ inputs.table_name }}. Zero anomalous
records detected. Ready for warehouse merge."
errors:
- id: alert_audit_error
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts on-call data engineers if the quality audit pipeline fails.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "⚠️ *Data Quality Flow Failure:* Execution {{ execution.id }} on flow {{ flow.id }} encountered an error."
}