id: csv-schema-contract-validator
namespace: company.team
description: |
Validate an inbound CSV against a JSON schema contract — required columns,
types, and row limits — and alert Slack when the contract breaks.
triggers:
- id: file_arrival_webhook
type: io.kestra.plugin.core.trigger.Webhook
description: Call from your upload or storage pipeline when a new CSV lands.
key: csv-contract-check
- id: daily_validation
type: io.kestra.plugin.core.trigger.Schedule
description: Optional daily re-validation of the same URL; shipped disabled.
cron: "0 7 * * *"
disabled: true
inputs:
- id: csv_url
type: STRING
defaults: https://example.com/data.csv
description: URL of the CSV to validate.
- id: schema_json
type: JSON
defaults: '{"required_columns": ["id", "name", "email"], "max_rows": 100000,
"non_empty_columns": ["id", "email"]}'
description: Contract — required columns, row limit, and columns that must never
be empty.
tasks:
- id: validate_csv
type: io.kestra.plugin.scripts.python.Script
description: Download the CSV, check required columns exist, enforce the row
limit, count empty cells in non-empty columns, and emit violations via the
stdout outputs protocol.
containerImage: python:3.11-slim
warningOnStdErr: false
env:
CSV_URL: "{{ inputs.csv_url }}"
SCHEMA: '{{ inputs.schema_json | toJson }}'
script: |
import csv, io, json, os, urllib.request
schema = json.loads(os.environ["SCHEMA"])
with urllib.request.urlopen(os.environ["CSV_URL"], timeout=60) as r:
text = r.read().decode("utf-8", "replace")
reader = csv.DictReader(io.StringIO(text))
rows = list(reader)
violations = []
missing = [c for c in schema.get("required_columns", []) if c not in (reader.fieldnames or [])]
if missing:
violations.append(f"missing columns: {missing}")
if len(rows) > schema.get("max_rows", float("inf")):
violations.append(f"row count {len(rows)} exceeds {schema['max_rows']}")
for col in schema.get("non_empty_columns", []):
if col in (reader.fieldnames or []):
empty = sum(1 for row in rows if not (row.get(col) or "").strip())
if empty:
violations.append(f"column '{col}' has {empty} empty rows")
print(f"rows: {len(rows)}, violations: {len(violations)}")
print("::" + json.dumps({"outputs": {"rows": len(rows), "violations": violations, "valid": len(violations) == 0}}) + "::")
- id: check_contract
type: io.kestra.plugin.core.flow.If
description: Valid → log the pass; violations → Slack with the actual violation list.
condition: "{{ outputs.validate_csv.vars.valid }}"
then:
- id: log_valid
type: io.kestra.plugin.core.log.Log
description: Record each passing validation so quiet weeks are visible too.
message: "CSV contract passed for {{ inputs.csv_url }}: {{
outputs.validate_csv.vars.rows }} rows validated."
else:
- id: alert_violations
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Post the violations list — not just "invalid" — so the producer can
fix the file.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": ":x: CSV contract BROKEN for {{ inputs.csv_url }} — {{ outputs.validate_csv.vars.violations }}. Execution {{ execution.id }}."
}
errors:
- id: alert_validation_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alert when validation itself errors — an unreadable file must not
pass as valid.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "CSV contract validation FAILED in flow {{ flow.id }} (execution {{ execution.id }}) for {{ inputs.csv_url }}. Check the URL and schema_json."
}
outputs:
- id: validation_result
type: JSON
description: 'Row count, violations, and validity flag, e.g. {"rows": 1204,
"violations": [], "valid": true}.'
value: "{{ outputs.validate_csv.vars | toJson }}"