New to Kestra?
Use blueprints to kickstart your first workflows.
Load a customer feed into a Postgres SCD Type 2 dimension, quarantine invalid rows, and alert Slack only when a version opens or closes.
id: postgres-customer-scd2
namespace: company.team
description: SCD Type 2 customer dimension in Postgres. Invalid rows are
quarantined; Slack fires only when history changes.
concurrency:
limit: 1
triggers:
- id: nightly
type: io.kestra.plugin.core.trigger.Schedule
description: Nightly merge at 02:00 UTC. Run manually first.
cron: "0 2 * * *"
timezone: UTC
inputs:
- id: snapshot
type: SELECT
displayName: Customer feed snapshot
description: "`baseline` is the first load. `address_change` updates Ada, adds
Katherine, and quarantines a row with no email."
defaults: baseline
values:
- baseline
- address_change
tasks:
- id: ensure_schema
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
description: Create staging, dimension, and quarantine tables if missing.
fetchType: NONE
sql: |
DO $scd$
BEGIN
CREATE TABLE IF NOT EXISTS stg_customer (
customer_id text NOT NULL,
full_name text NOT NULL,
email text,
city text,
plan text NOT NULL
);
CREATE TABLE IF NOT EXISTS dim_customer (
customer_key bigserial PRIMARY KEY,
customer_id text NOT NULL,
full_name text NOT NULL,
email text NOT NULL,
city text NOT NULL,
plan text NOT NULL,
valid_from timestamptz NOT NULL DEFAULT now(),
valid_to timestamptz,
is_current boolean NOT NULL DEFAULT true,
opened_by_run text,
closed_by_run text,
CONSTRAINT dim_customer_valid_window CHECK (valid_to IS NULL OR valid_to >= valid_from)
);
ALTER TABLE dim_customer ADD COLUMN IF NOT EXISTS opened_by_run text;
ALTER TABLE dim_customer ADD COLUMN IF NOT EXISTS closed_by_run text;
CREATE UNIQUE INDEX IF NOT EXISTS dim_customer_one_current
ON dim_customer (customer_id)
WHERE is_current;
CREATE TABLE IF NOT EXISTS quarantine_customer (
quarantine_id bigserial PRIMARY KEY,
customer_id text,
full_name text,
email text,
city text,
plan text,
reason text NOT NULL,
rejected_at timestamptz NOT NULL DEFAULT now(),
run_id text NOT NULL
);
END
$scd$;
- id: stage_feed
type: io.kestra.plugin.core.flow.If
description: Replace staging with the selected snapshot.
condition: "{{ inputs.snapshot == 'address_change' }}"
then:
- id: stage_address_change
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
description: Stage the changed feed.
fetchType: NONE
sql: |
WITH cleared AS (
DELETE FROM stg_customer
RETURNING 1
)
INSERT INTO stg_customer (customer_id, full_name, email, city, plan)
SELECT customer_id, full_name, email, city, plan
FROM (VALUES
('cust-100', 'Ada Lovelace', 'ada@example.com', 'Paris', 'pro'),
('cust-200', 'Grace Hopper', 'grace@example.com', 'Arlington', 'pro'),
('cust-300', 'Alan Turing', 'alan@example.com', 'Manchester', 'free'),
('cust-400', 'Katherine Johnson', 'katherine@example.com', 'Hampton', 'pro'),
('cust-500', 'Edsger Dijkstra', NULL, 'Utrecht', 'free')
) AS feed(customer_id, full_name, email, city, plan);
else:
- id: stage_baseline
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
description: Stage the original three customers.
fetchType: NONE
sql: |
WITH cleared AS (
DELETE FROM stg_customer
RETURNING 1
)
INSERT INTO stg_customer (customer_id, full_name, email, city, plan)
SELECT customer_id, full_name, email, city, plan
FROM (VALUES
('cust-100', 'Ada Lovelace', 'ada@example.com', 'Berlin', 'free'),
('cust-200', 'Grace Hopper', 'grace@example.com', 'Arlington', 'pro'),
('cust-300', 'Alan Turing', 'alan@example.com', 'Manchester', 'free')
) AS feed(customer_id, full_name, email, city, plan);
- id: quarantine_invalid_rows
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
description: Move rows with a missing email or city to quarantine.
fetchType: NONE
sql: |
WITH invalid AS (
DELETE FROM stg_customer
WHERE email IS NULL
OR btrim(email) = ''
OR city IS NULL
OR btrim(city) = ''
RETURNING
customer_id,
full_name,
email,
city,
plan,
CASE
WHEN email IS NULL OR btrim(email) = '' THEN 'missing_email'
ELSE 'missing_city'
END AS reason
),
fresh AS (
SELECT *
FROM invalid AS i
WHERE NOT EXISTS (
SELECT 1
FROM quarantine_customer AS q
WHERE q.customer_id IS NOT DISTINCT FROM i.customer_id
AND q.reason = i.reason
)
)
INSERT INTO quarantine_customer (customer_id, full_name, email, city, plan, reason, run_id)
SELECT customer_id, full_name, email, city, plan, reason, '{{ execution.id }}'
FROM fresh;
- id: close_changed_versions
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
description: Close the current row when name, email, city, or plan changed.
fetchType: NONE
sql: |
UPDATE dim_customer AS d
SET
is_current = false,
valid_to = now(),
closed_by_run = '{{ execution.id }}'
FROM stg_customer AS s
WHERE d.customer_id = s.customer_id
AND d.is_current
AND (
d.full_name IS DISTINCT FROM s.full_name
OR d.email IS DISTINCT FROM s.email
OR d.city IS DISTINCT FROM s.city
OR d.plan IS DISTINCT FROM s.plan
);
- id: insert_current_versions
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
description: Insert a current row for new keys and keys just closed.
fetchType: NONE
sql: |
INSERT INTO dim_customer (
customer_id,
full_name,
email,
city,
plan,
valid_from,
is_current,
opened_by_run
)
SELECT
s.customer_id,
s.full_name,
s.email,
s.city,
s.plan,
now(),
true,
'{{ execution.id }}'
FROM stg_customer AS s
WHERE s.email IS NOT NULL
AND btrim(s.email) <> ''
AND NOT EXISTS (
SELECT 1
FROM dim_customer AS d
WHERE d.customer_id = s.customer_id
AND d.is_current
);
- id: audit_dimension
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
description: Count opened, closed, and quarantined rows for this run.
fetchType: FETCH_ONE
sql: |
SELECT
(SELECT count(*) FROM quarantine_customer WHERE run_id = '{{ execution.id }}') AS quarantined,
(SELECT count(*) FROM dim_customer WHERE closed_by_run = '{{ execution.id }}') AS closed_versions,
(SELECT count(*) FROM dim_customer WHERE opened_by_run = '{{ execution.id }}') AS opened_versions,
(
SELECT count(*)
FROM (
SELECT customer_id
FROM dim_customer
GROUP BY customer_id
HAVING count(*) FILTER (WHERE is_current) <> 1
) AS broken
) AS broken_keys;
- id: assert_one_current_row
type: io.kestra.plugin.core.execution.Assert
description: Fail if any customer does not have exactly one current row.
errorMessage: "dim_customer has {{ outputs.audit_dimension.row.broken_keys }}
customer key(s) without exactly one current row."
conditions:
- "{{ outputs.audit_dimension.row.broken_keys == 0 }}"
- id: publish_changes
type: io.kestra.plugin.core.flow.If
description: Slack only when this run wrote history or quarantined a row.
condition: "{{ outputs.audit_dimension.row.quarantined > 0 or
outputs.audit_dimension.row.closed_versions > 0 or
outputs.audit_dimension.row.opened_versions > 0 }}"
then:
- id: notify_slack
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Post opened, closed, and quarantined counts.
url: "{{ secret('SLACK_WEBHOOK') }}"
messageText: "Customer dimension ({{ inputs.snapshot }}): opened {{
outputs.audit_dimension.row.opened_versions }}, closed {{
outputs.audit_dimension.row.closed_versions }}, quarantined {{
outputs.audit_dimension.row.quarantined }} (execution {{ execution.id
}})."
else:
- id: log_no_changes
type: io.kestra.plugin.core.log.Log
description: Log when the snapshot matches the current dimension.
message: "No dimension changes for snapshot {{ inputs.snapshot }}. Opened 0,
closed 0, quarantined 0 (execution {{ execution.id }})."
errors:
- id: alert_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Slack on flow failure.
url: "{{ secret('SLACK_WEBHOOK') }}"
messageText: "Customer dimension load failed in {{ flow.id }} (execution {{
execution.id }}). {{ errorLogs() }}"
Close the current customer row when name, email, city, or plan changes, insert a new current version, and quarantine rows with a missing email or city. A unique index plus an assertion keep exactly one is_current row per customer. Slack fires only when this run opened, closed, or quarantined a row; the same snapshot is a no-op. Overlapping merges are blocked with concurrency.limit: 1.
Prerequisites: Postgres (the flow creates stg_customer, dim_customer, and quarantine_customer) and a Slack incoming webhook.
Secrets:
POSTGRES_URL: JDBC URL (for example jdbc:postgresql://postgres:5432/warehouse)POSTGRES_USER / POSTGRES_PASSWORD: role that can create tablesSLACK_WEBHOOK: incoming webhook URLInput snapshot: baseline (default) or address_change.
Quick start: set secrets, run baseline (3 opened), run it again (no Slack), then address_change (2 opened, 1 closed, 1 quarantined). Counts are {{ outputs.audit_dimension.row.opened_versions }}, .closed_versions, .quarantined, and .broken_keys.
Links: PostgreSQL, If, Slack. Created by parshipcy.
