Schedule icon
If icon
Queries icon
Query icon
OutputValues icon
Log icon
Fail icon
Commands icon
Docker icon
Set icon
LoopUntil icon
Switch icon
UploadFiles icon

Upgrade Postgres to a new major version with logical replication and a verified cutover

Upgrade Postgres across major versions with logical replication in Kestra. Preflight, checksums, frozen cutover, sequence sync and automatic rollback.

Categories
DataInfrastructure

pg_upgrade needs the old and the new binaries on one host and a maintenance window that grows with the cluster. Logical replication avoids both: the new cluster subscribes to the old one, catches up while the application keeps writing, and the only downtime is the few seconds it takes to freeze the old cluster, drain the last changes and switch. It is also easy to get wrong in ways that only show up on cutover night: a table without a primary key breaks writes on production, sequences are not replicated, and a pooled connection opened during the freeze stays read-only afterwards.

This flow runs the whole upgrade as one auditable pipeline and handles each of those cases. Run it in REHEARSAL mode as often as you like: it proves the path, measures it and cleans up. Run it once in CUTOVER mode to switch.

This blueprint was created by zkasuran.

The phases

Phase What happens Reversible
1. Preflight Reads both clusters and lists every blocker with its fix. Nothing is written yes, nothing changed
2. Setup Fixes keyless tables if asked, copies the schema, creates the publication and the subscription yes, errors removes it
3. Catch-up Waits until every table finished its initial copy and the slot lag is under max_lag_bytes yes
4. Verification Row count and md5 over every row, per table, at zero lag yes
5a. Rehearsal end Removes the replication, optionally empties the target, writes the report yes
5b. Cutover Freezes the source, drains, copies sequences, detaches, smoke-tests the target until the application moves, errors unfreezes the source
6. Evidence Markdown report with every check, stored in namespace files

What the preflight blocks on

Every blocker names its fix. The run stops before touching either cluster.

  • wal_level is not logical on the source.
  • No free replication slot or WAL sender on the source.
  • The user lacks the REPLICATION attribute.
  • The target is not on a newer major version.
  • The target already holds tables or an old subscription (unless the run is resuming a rolled-back cutover, see below).
  • Tables without a primary key or replica identity. Once a table is in a publication, every UPDATE and DELETE on it fails on the source with cannot update table because it does not have a replica identity. That is a production outage caused by the upgrade preparation itself. Add a key, or set fix_replica_identity to apply REPLICA IDENTITY FULL.
  • Unlogged tables, which logical replication skips.
  • Extensions used on the source that the target cannot install.
  • CUTOVER without confirm_database matching the database name.
  • An earlier cutover in state CUTOVER_STARTED or CUTOVER_DONE.

Warnings are reported, not blocking: large objects (not replicated), the WAL cost of REPLICA IDENTITY FULL, leftover publications and slots, and the reminder that DDL is not replicated.

How it works

  1. demo (io.kestra.plugin.core.flow.If) creates the demo schema when setup_demo is on: 20,000 customers, 100,000 orders and about 250,000 order items with foreign keys, a composite key, a partial index and two sequences. demo_table_without_key adds a keyless audit_events table.
  2. source_facts and target_facts (io.kestra.plugin.jdbc.postgresql.Query) read versions, settings, keyless and unlogged tables, large objects, extensions, sizes and leftovers.
  3. state (io.kestra.plugin.core.output.OutputValues) reads the phase recorded in the KV store by earlier runs. preflight and preflight_extensions build the blocker and warning lists. preflight_gate (io.kestra.plugin.core.flow.If) stops on any blocker.
  4. setup_replication (io.kestra.plugin.core.flow.If), skipped when resuming:
    • set_replica_identity_full (io.kestra.plugin.scripts.shell.Commands, postgres:17 image) when fix_replica_identity is on.
    • record_started (io.kestra.plugin.core.kv.Set) records phase REPLICATING.
    • copy_schema: pg_dump --schema-only from the newer client, applied to the target with psql. Data is left to replication.
    • create_publication (FOR ALL TABLES) on the source and create_subscription on the target. CREATE SUBSCRIPTION cannot run inside a transaction, which is why these use psql and not a JDBC task.
    • writes_during_copy (demo) inserts, updates and deletes rows on the source while the copy runs.
  5. wait_for_sync (io.kestra.plugin.core.flow.LoopUntil) checks every 10 seconds, for up to 30 minutes, that pg_subscription_rel shows every table in state r and that the slot is within max_lag_bytes.
  6. verify runs at zero lag, in UTC, and computes for each published table count(*) and md5(string_agg(row::text, '|' ORDER BY row::text)) on both clusters. verify_gate fails on any mismatch and leaves the replication in place for inspection.
  7. by_mode (io.kestra.plugin.core.flow.Switch):
    • REHEARSAL: teardown_rehearsal drops the subscription (and with it the slot on the source) and the publication, and with rehearsal_cleanup empties the target. record_rehearsal stores the counts and timing.
    • CUTOVER:
      • record_cutover_started writes phase CUTOVER_STARTED.
      • freeze_source sets default_transaction_read_only = on on the source database and ends the open client sessions.
      • drain (LoopUntil) waits for zero lag.
      • sync_sequences copies every sequence value, because logical replication does not carry them.
      • detach drops the subscription and the publication.
      • smoke_test checks, for every sequence-backed column on the target, that the sequence is at or above the column's maximum. A missed sequence would fail the first insert with a duplicate key.
      • record_cutover_done writes phase CUTOVER_DONE.
  8. report and publish_report write pg-upgrade/<database>/<timestamp>-<mode>.md.
  9. errors reads the recorded phase once (failed_phase), then:
    • CUTOVER_STARTED: unfreeze_source makes the source writable again, and record_rolled_back keeps phase REPLICATING. The subscription and the target data stay.
    • REPLICATING: remove_replication drops the subscription, the publication and the slot, empties the target and resets the phase.
    • Anything else: nothing to undo.

Retrying a cutover that rolled back

When a cutover fails after the freeze, the errors branch unfreezes the source and keeps the subscription running, so the target keeps receiving writes. Start the cutover again with the same inputs. The preflight sees phase REPLICATING and a live subscription, setup_replication is skipped, and the run continues at the sync check. In testing, a row written to the source after the rollback was on the target when the retried cutover completed.

Tested end to end

On Kestra 2.0.3 OSS with Postgres 16.15 as the source and 17.11 as the target, using the demo schema (370,923 rows).

Run Inputs Result
Rehearsal setup_demo, demo_write_load 3 tables in sync after 2 checks, every count and md5 matches, including the 500 inserts, updates and deletes made during the copy. Replication removed, target emptied
Rehearsal again defaults Same result. Rehearsals are repeatable
Keyless table demo_table_without_key Blocked in preflight: no primary key or replica identity on public.audit_events. Nothing changed
Keyless table fixed fix_replica_identity REPLICA IDENTITY FULL applied, 4 tables verified, audit_events md5 matches
Cutover without confirmation mode: CUTOVER Blocked in preflight: CUTOVER needs confirm_database set to shop
Cutover confirm_database: shop, demo_write_load Source frozen, drained, 2 sequences copied, smoke test passed. Afterwards the source refuses writes, the target has no subscription, and the first new order on the target got id 100501
Second cutover same Blocked: phase CUTOVER_DONE and a non-empty target
Failure after the freeze demo_fail_after_freeze Source writable again, subscription and target data kept, phase back to REPLICATING
Retry after the rollback same as the cutover Resumed without a new copy, completed, and the row written during the rollback reached the target

Problems found while testing

  • Pooled connections keep the freeze. Kestra JDBC tasks pool connections. A connection opened while the source was frozen stayed read-only after the freeze was lifted, so the next run's first statement failed with cannot execute DROP TABLE in a read-only transaction. Every task that talks to the source sets connectionPooling: false.
  • Undoing the freeze needs an opt-out. Once default_transaction_read_only is on, the ALTER DATABASE ... RESET that undoes it is itself refused. The undo and the detach run with PGOPTIONS=-c default_transaction_read_only=off for that one session.
  • A Switch on a KV value re-renders. The first errors branch switched directly on the KV entry, and the undo steps change that entry, so more than one branch ran. The phase is now read once into failed_phase.
  • setval needs a quoted boolean. format('%s', true) renders t, which setval reads as a column name. The sequence sync uses %L.
  • srsubstate is a "char". Concatenating it with text is ambiguous until it is cast to text.
  • Typed confirmation is an input, not a pause. A Pause before the freeze looked natural, but on 2.0.3 resuming that Pause was unreliable in this flow, and a run paused while holding the concurrency slot blocks every other run of the flow. The confirmation is therefore an input checked in preflight. Read the rehearsal report, then start the cutover deliberately.

Inputs

Input Type Default Purpose
mode SELECT REHEARSAL REHEARSAL or CUTOVER
source_host, source_port STRING, INT pg-source, 5432 The old cluster, reachable from the worker and from the target
target_host, target_port STRING, INT pg-target, 5432 The new cluster
database STRING shop Must exist, empty, on the target
username STRING postgres Replication user, the same on both clusters
max_lag_bytes INT 65536 In-sync threshold
fix_replica_identity BOOL false Apply REPLICA IDENTITY FULL to keyless tables
confirm_database STRING empty Must equal database for CUTOVER
rehearsal_cleanup BOOL true Empty the target after a rehearsal
setup_demo, demo_table_without_key, demo_write_load, demo_fail_after_freeze BOOL false Demo only

Prerequisites

  • Source: wal_level = logical (a restart), enough max_replication_slots and max_wal_senders, and a pg_hba.conf entry that lets the target connect for replication.
  • Target: the new major version installed, the database created empty, the roles that own objects created, and the source's extension packages installed.
  • Both clusters reachable from the Kestra worker, and the source reachable from the target.
  • A Kestra worker that can run Docker containers, for the postgres:17 client image. Use the client image of the target version.

Secrets

  • PG_SOURCE_PASSWORD: password of username on the source. It is also embedded in the subscription connection string, which is stored in pg_subscription on the target until the subscription is dropped.
  • PG_TARGET_PASSWORD: password of username on the target.

Runbook

  1. Weeks before: enable wal_level = logical on the source during a normal maintenance window. Install the new cluster.
  2. Run a REHEARSAL on a copy of production, then on production itself. Fix every blocker. Read the report: the number of sync checks times 10 seconds is your catch-up time.
  3. Enable the weekly_rehearsal schedule until cutover day, so drift such as a new keyless table is caught early.
  4. Cutover day: freeze schema changes, announce the window, run CUTOVER with confirm_database.
  5. When the run succeeds, point the application at the target, using the connection string in the report. Keep the old cluster read-only for a few days as a fallback, then retire it.
  6. If the run fails after the freeze, the source is writable again automatically. Read the failed task, fix the cause, and run the cutover again. It resumes from the live subscription.

Expected outputs

  • outputs.verify.vars: tables, rows, mismatches, lag_bytes, and verification.tsv with per-table counts and hashes.
  • Namespace file pg-upgrade/<database>/<timestamp>-<mode>.md.
  • KV pg_upgrade_<database>: the phase and the last run's facts.

Things to know

  • DDL is not replicated. Any schema change between the subscription and the cutover must be applied to both clusters by hand, or the rehearsal must be redone.
  • The md5 check reads every row of every table on both clusters. On very large tables, run it on a sample or a recent partition.
  • Sequences, large objects, materialized views and their data are not replicated. Sequences are handled. Refresh materialized views on the target after the cutover.
  • The freeze applies to new sessions and ends existing client sessions. Connection poolers in front of the source keep reconnecting and get read-only sessions, which surfaces as write errors in the application until it moves.

How to extend

  • Gate the cutover on a change ticket by reading its status with an HTTP task before freeze_source.
  • Switch the application automatically, for example by updating a DNS record, a PgBouncer config or a Kubernetes secret, right after record_cutover_done.
  • Replace FOR ALL TABLES with a table list to move a database in stages.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.