OutputValues icon
Loop icon
Script icon
UploadFiles icon
Set icon
If icon
Log icon
SlackIncomingWebhook icon
Schedule icon

Incremental Postgres to S3 Parquet replication with Sling and Kestra watermarks

Replicate changed Postgres rows to S3 Parquet with Sling. Kestra keeps per-table watermarks, runs tables in parallel and alerts on volume spikes.

Categories
Data

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:

  • copies the rows whose update timestamp falls in (last watermark, run start - 1 minute] to a new Parquet file per table in S3, all tables in parallel,
  • advances a table's watermark only once its file is written, so a failed table is retried from the same point by the next run,
  • warns, and alerts Slack when 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.

How it works

  1. 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.
  2. 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.
  3. 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.
    • In demo mode, 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.
  4. If a run fails, alert_failed_run (io.kestra.plugin.slack.notifications.SlackIncomingWebhook) posts to Slack. Tables that succeeded in that run keep their new watermark.

Inputs

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

Prerequisites

  • A Kestra worker that can run Docker containers. The script uses python:3.12-slim and installs the pinned Sling version (and DuckDB in demo mode) for each table, which takes about 15 seconds.
  • For production: Postgres tables with an update timestamp column that is set on every insert and update (for example with a trigger or your ORM), and an S3 bucket.

Secrets

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.

Quick start

  1. Run the flow. Each table copies everything since midnight UTC (up to about 8,600 orders and 3,400 payments late in the day; close to midnight there is little to copy). The files appear in the namespace Files tab under sling_demo/.
  2. Wait a minute and run it again. Only the rows that arrived in between are copied, for example 6 orders and 2 payments. The watermarks are in the KV Store under sling_demo_watermark_*.
  3. For production, set the secrets, change bucket, region and prefix in variables and the default tables, and set disabled: false on the hourly trigger.

Expected outputs

  • Per table: outputs.sync.vars.rows, low, high, initial_load, baseline and spike, and a Parquet file per non-empty batch.
  • KV Store: sling_watermark_<schema>__<table> (the last high) and sling_rows_<schema>__<table> (the last batch size), with a sling_demo_ prefix in demo mode.
  • A WARN log line, and a Slack message when notify_slack is on, when a batch size changes by more than spike_factor.

Things to know

  • Rows are captured by their update timestamp, so hard deletes are not replicated. Use soft deletes, or a CDC tool such as Debezium.
  • Rows are missed if a transaction stays open longer than 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.
  • A row updated twice between two runs is copied once, with its latest values. Downstream, keep the latest version of each primary key.
  • Volume checks start once a table has an incremental batch of at least 10 rows. After a spike, the next normal batch can raise a second alert in the other direction.
  • Table names must be unquoted schema.table identifiers. To re-copy a table from scratch, delete its watermark from the KV Store.

How to extend

  • Load the Parquet files into a warehouse with its own loader, for example Snowflake 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.
  • Add tables by appending them to tables; each one gets its own watermark.
  • Replace the Schedule with a Postgres trigger to run as soon as new rows appear.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.