New to Kestra?
Use blueprints to kickstart your first workflows.
Validate S3 files with Great Expectations in Kestra. Load passing files to Postgres, quarantine failures and alert Slack. Runs with no setup.
Bad files should never reach the warehouse. This flow puts a Great Expectations (GX) gate between an S3 landing prefix and Postgres. Every new file is checked against an expectation suite, in parallel:
The expectation suite is plain JSON inside the flow, so you can change the rules without touching the Python code.
It runs with no setup. On a manual run the flow validates a public CSV of 100 orders and a copy of it with four bad rows, so you can watch one file pass and the other get rejected. S3 copies only happen on S3-triggered runs; loading and Slack also run on manual runs when you switch them on with inputs.
This blueprint was created by manojgosavi.
new_files (io.kestra.plugin.aws.s3.Trigger) polls the landing prefix every minute for .csv files and starts an execution with up to 20 of them. The files are downloaded to Kestra's internal storage and deleted from the landing prefix (action: DELETE), so a file is never picked up twice.download_sample (io.kestra.plugin.core.http.Download) fetches the sample CSV instead.plan_batch (io.kestra.plugin.core.output.OutputValues) builds the list of files (the S3 objects, or the two demo files) and records whether the run came from S3.validate_files (io.kestra.plugin.core.flow.Loop) processes up to 4 files at a time. For each file:validate (io.kestra.plugin.scripts.python.Script, great_expectations==1.23.2) runs the suite with an ephemeral GX context and outputs passed, rows, evaluated, failed_count, failed and a one-line summary. A file that can't be read as CSV fails with the summary unreadable CSV: <error>.route_file (io.kestra.plugin.core.flow.If):load_to_postgres (io.kestra.plugin.jdbc.postgresql.CopyIn) appends the rows to public.orders and archive_file (io.kestra.plugin.aws.s3.Upload) writes the file to processed/orders/<date>/<execution id>_<file name>.log_failure (io.kestra.plugin.core.log.Log) logs the reasons, quarantine_file (io.kestra.plugin.aws.s3.Upload) writes the file to quarantine/orders/<date>/<execution id>_<file name> and alert_slack (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) posts the failing expectations.alert_failed_run (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) posts to Slack so the files are not silently stuck in internal storage.The suite.json input file of validate checks that:
order_id, customer_name, customer_email, product_id, price, quantity, total, in this order (CopyIn loads columns by position),order_id is never null and is unique,customer_email looks like an email address,price is never null and is greater than 0,quantity is never null and is between 1 and 100.Each entry is an expectation type and its kwargs, as listed in the Expectation Gallery. Add, remove or tune entries to match your data.
sample_url (URI, default the public orders CSV): the file validated on manual runs.inject_bad_rows (BOOL, default true): also validate orders_with_bad_rows.csv, a copy of the sample with a duplicate id, an invalid email, a negative price and a quantity of 500.load_to_postgres (BOOL, default false): load passing files on manual runs too.notify_slack (BOOL, default false): send Slack alerts on manual runs too.python:3.12-slim and installs the pinned GX version for each file, which takes about 20 seconds. For high volumes, build an image with GX preinstalled and drop beforeCommands.CREATE TABLE public.orders (order_id int, customer_name text, customer_email text, product_id int, price numeric, quantity int, total numeric);
Only needed for S3-triggered runs, or when you enable loading or alerts on manual runs:
AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY: read and delete objects under the landing prefix, write under the processed and quarantine prefixes.POSTGRES_URL (for example jdbc:postgresql://host:5432/warehouse), POSTGRES_USER, POSTGRES_PASSWORD: insert into the target table.SLACK_WEBHOOK_URL: Slack incoming webhook for quarantine alerts.orders.csv passes all 8 expectations and orders_with_bad_rows.csv fails 4 of them and is rejected with a WARN log line. On S3-triggered runs it would also be copied to quarantine and reported to Slack.inject_bad_rows: false. Only orders.csv is validated, and it passes.bucket, region, the prefixes and target_table in variables, and set disabled: false on the new_files trigger. Then drop a CSV under landing/orders/.outputs.validate.vars.passed, rows, evaluated, failed_count, failed (for example expect_column_values_to_be_unique on order_id (2 rows)) and summary.public.orders and a copy under processed/orders/<date>/.WARN log line and, on S3-triggered runs, a copy under quarantine/orders/<date>/ and a Slack message.alert_failed_run tells you, and restarting the execution from the Executions page processes them again.CopyIn with a load into Snowflake, BigQuery or Databricks.