New to Kestra?
Use blueprints to kickstart your first workflows.
Replicate changed Postgres rows to S3 Parquet with Sling. Kestra keeps per-table watermarks, runs tables in parallel and alerts on volume spikes.
Sling moves data between databases and object storage with a single command. The free CLI does not keep incremental state when the target is a file, so every run would copy whole tables again. This flow lets Kestra own the state: it keeps a watermark per table in the KV Store and asks Sling only for the rows changed since then.
Each run:
(last watermark, run start - 1 minute] to a new Parquet file per table in S3, all tables in parallel,notify_slack is on, when a batch is more than spike_factor (default 5) times larger or smaller than the previous non-empty batch.It runs with no setup. In demo mode the source is a DuckDB database generated inside the task, where an order arrives every 10 seconds and a payment every 25 seconds since midnight UTC, and the Parquet files go to the namespace files. Run it twice a minute apart and the second run copies only the rows that arrived in between.
This blueprint was created by manojgosavi.
hourly (io.kestra.plugin.core.trigger.Schedule) starts the flow every hour with demo_mode: false and notify_slack: true. Runs are queued one at a time (concurrency.limit: 1), and after downtime only the last missed run is started, which covers the whole gap.plan (io.kestra.plugin.core.output.OutputValues) fixes the upper bound of the window (high, the execution start minus lag_minutes) and picks the tables: inputs.tables, or main.orders and main.payments in demo mode.replicate (io.kestra.plugin.core.flow.Loop) processes up to 3 tables at a time. For each table:sync (io.kestra.plugin.scripts.python.Script, sling==1.6.4) reads the table's watermark from the KV Store and runs Sling with SELECT * FROM <table> WHERE <update_key> > <watermark> AND <update_key> <= <high>. The batch goes to raw/postgres/<schema>/<table>/<execution start, UTC, yyyyMMdd_HHmmss>.parquet in S3. It outputs rows, low, high, initial_load, baseline and spike.publish_demo_files (io.kestra.plugin.core.namespace.UploadFiles) copies the batch to the namespace files under sling_demo/.save_watermark (io.kestra.plugin.core.kv.Set) stores high as the new watermark.check_volume (io.kestra.plugin.core.flow.If): when spike is true, log_volume_change (io.kestra.plugin.core.log.Log) logs a warning and alert_slack (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) posts it.save_row_count (io.kestra.plugin.core.kv.Set) keeps the batch size as the next baseline. Empty batches and the initial full load are not used as a baseline.alert_failed_run (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) posts to Slack. Tables that succeeded in that run keep their new watermark.demo_mode (BOOL, default true): use the generated DuckDB source and namespace files. The trigger sets it to false.tables (ARRAY of STRING, default public.customers, public.payments): the Postgres tables, as schema.table.update_key (STRING, default updated_at): the timestamp column set on every insert and update.spike_factor (INT, default 5): the volume alert threshold.notify_slack (BOOL, default false): post volume alerts and failed runs to Slack. The trigger sets it to true.python:3.12-slim and installs the pinned Sling version (and DuckDB in demo mode) for each table, which takes about 15 seconds.Only needed when demo_mode is false:
POSTGRES_HOST, POSTGRES_DB, POSTGRES_USER, POSTGRES_PASSWORD: read access to the tables. The SSL mode is set in variables.pg_sslmode (require; use verify-full to also verify the server certificate).AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY: write access to the bucket.SLACK_WEBHOOK_URL: Slack incoming webhook for volume alerts and failed runs.sling_demo/.sling_demo_watermark_*.bucket, region and prefix in variables and the default tables, and set disabled: false on the hourly trigger.outputs.sync.vars.rows, low, high, initial_load, baseline and spike, and a Parquet file per non-empty batch.sling_watermark_<schema>__<table> (the last high) and sling_rows_<schema>__<table> (the last batch size), with a sling_demo_ prefix in demo mode.WARN log line, and a Slack message when notify_slack is on, when a batch size changes by more than spike_factor.lag_minutes, or if the Postgres clock runs behind Kestra's. Raise lag_minutes to cover both.update_key should be timestamptz. With timestamp without time zone, the database TimeZone must match how the timestamps were written.schema.table identifiers. To re-copy a table from scratch, delete its watermark from the KV Store.COPY INTO or BigQuery external tables. Writing to a warehouse table directly with Sling also needs --mode incremental and --primary-key, because full-refresh would replace the table with each batch.tables; each one gets its own watermark.