Download icon
OutputValues icon
Loop icon
Script icon
If icon
CopyIn icon
Upload icon
Log icon
SlackIncomingWebhook icon
Trigger icon

Quarantine bad files from S3 with a Great Expectations data quality gate

Validate S3 files with Great Expectations in Kestra. Load passing files to Postgres, quarantine failures and alert Slack. Runs with no setup.

Categories
Data

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:

  • Files that pass are appended to the target table and archived under a processed prefix.
  • Files that fail are moved to a quarantine prefix, are never loaded, and Slack gets the list of failing expectations with the number of offending rows.

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.

How it works

  1. 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.
  2. On manual runs, download_sample (io.kestra.plugin.core.http.Download) fetches the sample CSV instead.
  3. 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.
  4. 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):
      • Passed: 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>.
      • Failed: 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.
  5. If the execution itself fails, alert_failed_run (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) posts to Slack so the files are not silently stuck in internal storage.

The expectation suite

The suite.json input file of validate checks that:

  • the columns are exactly 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.

Inputs

  • 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.

Prerequisites

  • A Kestra worker that can run Docker containers. The script uses 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.
  • For S3-triggered runs: an S3 bucket, and a Postgres table whose columns match the CSV header, for example:
    CREATE TABLE public.orders (order_id int, customer_name text, customer_email text, product_id int, price numeric, quantity int, total numeric);
    

Secrets

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.

Quick start

  1. Run the flow. In the Logs, 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.
  2. Run it again with inject_bad_rows: false. Only orders.csv is validated, and it passes.
  3. For production, create the table, set the secrets, change 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/.

Expected outputs

  • Per file: outputs.validate.vars.passed, rows, evaluated, failed_count, failed (for example expect_column_values_to_be_unique on order_id (2 rows)) and summary.
  • Passing files: rows in public.orders and a copy under processed/orders/<date>/.
  • Failing files: a WARN log line and, on S3-triggered runs, a copy under quarantine/orders/<date>/ and a Slack message.

Things to know

  • A file is either loaded in full or not at all. The gate does not split good rows from bad ones.
  • The trigger deletes files from the landing prefix as soon as they are downloaded, before validation. If an execution fails (for example a network error while installing GX, or a Postgres outage), the files only exist in Kestra's internal storage: alert_failed_run tells you, and restarting the execution from the Executions page processes them again.
  • Archived and quarantined copies are prefixed with the date and execution id, so files with the same name never overwrite each other.
  • To re-process a quarantined file after fixing it, copy it back under the landing prefix.

How to extend

  • Keep one suite per dataset in namespace files and pick it from the file's prefix.
  • Build GX Data Docs from the validation result and publish them as namespace files.
  • Replace CopyIn with a load into Snowflake, BigQuery or Databricks.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.