id: s3-csv-duckdb-slack-quality-alert
namespace: company.data
inputs:
- id: s3_bucket
displayName: S3 Bucket Name
type: STRING
defaults: my-data-warehouse-landing
description: Target S3 bucket where raw CSV files are ingested.
- id: s3_prefix
displayName: S3 Key Prefix
type: STRING
defaults: raw/transactions/
description: S3 directory prefix to poll for new data batches.
- id: quality_threshold
displayName: Data Quality Completeness Threshold
type: FLOAT
defaults: 0.95
description: Minimum required non-null ratio across all columns (0.0 to 1.0).
triggers:
- id: daily_audit
type: io.kestra.plugin.core.trigger.Schedule
cron: "0 6 * * *"
description: Trigger automated data quality checks daily at 06:00 UTC.
tasks:
- id: list_s3_files
type: io.kestra.plugin.aws.s3.List
accessKeyId: "{{ secret('AWS_ACCESS_KEY_ID') }}"
secretKeyId: "{{ secret('AWS_SECRET_ACCESS_KEY') }}"
region: us-east-1
bucket: "{{ inputs.s3_bucket }}"
prefix: "{{ inputs.s3_prefix }}"
description: Scans the target S3 bucket prefix to find incoming files.
- id: download_latest_file
type: io.kestra.plugin.aws.s3.Download
accessKeyId: "{{ secret('AWS_ACCESS_KEY_ID') }}"
secretKeyId: "{{ secret('AWS_SECRET_ACCESS_KEY') }}"
region: us-east-1
bucket: "{{ inputs.s3_bucket }}"
key: "{{ inputs.s3_prefix }}latest.csv"
description: Downloads the latest batch CSV into Kestra's execution context.
- id: duckdb_quality_engine
type: io.kestra.plugin.scripts.python.Script
description: Spins up DuckDB to execute column-level schema and completeness validation.
docker:
image: python:3.12-slim
beforeCommands:
- pip install duckdb pandas
inputFiles:
input_data.csv: "{{ outputs.download_latest_file.uri }}"
script: |
import duckdb
import json
import sys
conn = duckdb.connect(':memory:')
conn.execute("CREATE TABLE batch_data AS SELECT * FROM read_csv_auto('input_data.csv')")
total_rows = conn.execute("SELECT COUNT(*) FROM batch_data").fetchone()[0]
columns_info = conn.execute("PRAGMA table_info('batch_data')").fetchall()
report = {
"total_rows": total_rows,
"column_metrics": {},
"null_anomalies": []
}
if total_rows == 0:
report["overall_completeness"] = 0.0
report["null_anomalies"].append("Empty dataset received")
else:
scores = []
for col in columns_info:
col_name = col[1]
null_count = conn.execute(f"SELECT COUNT(*) FROM batch_data WHERE \"{col_name}\" IS NULL").fetchone()[0]
ratio = 1.0 - (null_count / total_rows)
report["column_metrics"][col_name] = {
"null_count": null_count,
"completeness_score": round(ratio, 4)
}
scores.append(ratio)
if ratio < 0.90:
report["null_anomalies"].append(f"Column '{col_name}' has excessive nulls: {null_count}/{total_rows}")
report["overall_completeness"] = round(sum(scores) / len(scores), 4) if scores else 0.0
with open("{{ outputDir }}/quality_summary.json", "w") as f:
json.dump(report, f, indent=2)
with open("{{ outputDir }}/score.txt", "w") as f:
f.write(str(report["overall_completeness"]))
print(f"Audit completed: {total_rows} rows analyzed. Overall score: {report['overall_completeness']}")
outputFiles:
- quality_summary.json
- score.txt
- id: evaluate_quality_gate
type: io.kestra.plugin.core.flow.If
description: Evaluates whether the calculated score satisfies the SLA threshold.
condition: "{{ outputs.duckdb_quality_engine.vars.score < inputs.quality_threshold }}"
then:
- id: notify_slack_degradation
type: io.kestra.plugin.notifications.slack.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "🚨 Data Quality SLA Breached",
"blocks": [
{
"type": "header",
"text": {"type": "plain_text", "text": "🚨 Data Quality Alert: SLA Breached"}
},
{
"type": "section",
"fields": [
{"type": "mrkdwn", "text": "*Calculated Score:* {{ outputs.duckdb_quality_engine.vars.score }}"},
{"type": "mrkdwn", "text": "*Required Threshold:* {{ inputs.quality_threshold }}"},
{"type": "mrkdwn", "text": "*Target:* `s3://{{ inputs.s3_bucket }}/{{ inputs.s3_prefix }}`"},
{"type": "mrkdwn", "text": "*Execution ID:* `{{ execution.id }}`"}
]
}
]
}
description: Alerts engineering team on Slack when data fails validation.
else:
- id: notify_slack_health
type: io.kestra.plugin.notifications.slack.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "✅ Data Quality Verified: Score {{ outputs.duckdb_quality_engine.vars.score }} satisfies threshold {{ inputs.quality_threshold }}"
}
description: Confirms data quality passed without incident.