Trigger icon
Schedule icon
Query icon
If icon
Fail icon
SlackIncomingWebhook icon
Pause icon
Log icon

Stop Stale Replication Slots from Filling Your Postgres Disk

Detect inactive PostgreSQL replication slots retaining WAL, alert without spam, and drop dead slots safely after on-call approval.

Categories
DataInfrastructure

A replication slot makes PostgreSQL keep every WAL segment its consumer has not confirmed yet. When the consumer stops (a Debezium connector that crashed, a CDC job that was switched off, a replica that was deleted), the slot keeps holding WAL until the disk fills and the database stops accepting writes. This blueprint watches every slot, tells you early and only when something changes, and helps you clean up dead slots without ever dropping one whose consumer is still alive.

How it works

  1. wal_retention (io.kestra.plugin.jdbc.postgresql.Trigger) polls every minute and starts a run only when a slot holds more WAL than the warning threshold or has lost its reserved WAL. hourly_check (Schedule) runs regardless, to report recoveries.
  2. check_server fails fast if the database is unreachable, and records whether the role may drop slots. version_guard (If) requires PostgreSQL 13+.
  3. inspect_slots measures the WAL each slot retains (pg_wal_lsn_diff from restart_lsn) and classifies it in one SQL statement:
    • ok, warn (warn_mb), critical (critical_mb, or wal_status unreserved/lost)
    • it updates the replication_slot_watch table with the level, position, and when the slot became inactive
    • slots that disappeared are reported as removed
  4. report_changes (If) posts only the slots whose level changed, or that just lost their WAL, with what to do about each.
  5. propose_drops (If) lists critical slots inactive for drop_after_hours. If the role may drop slots, wait_for_decision (Pause) waits up to 4 hours for the on-call engineer to enter the names to drop.
  6. drop_slots drops a slot only if it was proposed, is still inactive, and its position has not moved since the proposal. Every name gets a result line in Slack.

Error handling

  • No alert storms. A slot is reported when its level changes, so a slot that stays critical does not page every minute. A proposal is not repeated for ask_again_hours.
  • No unsafe drops. The drop is re-checked in the same statement that performs it. A slot whose consumer reconnected, or that advanced since the proposal, is refused with the reason, for example "refused: it advanced since the proposal (0/28E3B080 to 0/35C4FDF0)". Names that were never proposed (typos, healthy slots) are refused too.
  • No answer means keep. The approval Pause resumes after 4 hours with nothing selected.
  • Least privilege. Monitoring only needs to read pg_replication_slots. When the role lacks REPLICATION, the proposal says so and the flow does not wait for an approval it cannot carry out.
  • Clear failures. check_server retries with backoff; if the database stays unreachable, the errors block names the task and cause and says that nothing was dropped.

What this blueprint teaches about Kestra

  • An event-based JDBC Trigger whose query is the alert condition, so the flow runs only when there is something to look at.
  • Keeping state between runs in a table to turn level-triggered checks into change notifications.
  • An approval Pause with behavior: RESUME as a safe default, and a re-check at execution time, so human approval never acts on stale data.
  • Bound SQL parameters for values typed by people.

Prerequisites

  • PostgreSQL 13 or newer.
  • A role that can read pg_replication_slots and create the replication_slot_watch table, for example CREATE ROLE slot_guard LOGIN PASSWORD '...'; GRANT pg_monitor TO slot_guard; GRANT CREATE ON SCHEMA public TO slot_guard;. Add REPLICATION (ALTER ROLE slot_guard REPLICATION;) if the flow should drop slots.
  • A Slack incoming webhook.

Secrets

  • POSTGRES_URL, POSTGRES_USER, POSTGRES_PASSWORD: JDBC URL of the database to watch (for example jdbc:postgresql://db:5432/app) and the role above.
  • SLACK_WEBHOOK_URL: Slack incoming webhook.

Inputs

  • warn_mb (INT, default 1024), critical_mb (INT, default 4096): WAL retained per slot, in MiB. Keep the trigger's threshold in sync with warn_mb.
  • drop_after_hours (INT, default 24): how long a critical slot must be inactive before it is proposed for dropping.
  • ask_again_hours (INT, default 12): how long to wait before proposing the same slot again.
  • kestra_url (STRING): base URL for the approval link.

Outputs

  • slots (JSON): every slot with its retained WAL in MiB, level, activity, and whether it was proposed for dropping.

Quick start

  1. Create the role and add the secrets.
  2. Run the flow manually once: it creates replication_slot_watch and reports any slot already above the thresholds.
  3. Enable wal_retention and hourly_check.
  4. Optionally set max_slot_wal_keep_size in PostgreSQL as a hard limit: slots that pass it become lost instead of filling the disk, and this flow reports them.

How to extend

  • Page through PagerDuty or Opsgenie instead of Slack for critical changes.
  • Watch several databases by turning check_server and inspect_slots into a subflow called once per database.
  • Restart the consumer automatically, for example a Debezium connector through the Kafka Connect REST API, before proposing a drop.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.