Consolidating databases with Postgres logical replication
This is part 2 of a series about Postgres logical replication use-cases, and about how the feature set has evolved over the past ten years and ten releases, one architecture at a time. Part 1 built a hub-and-workers system for write scalability and has the table of what each release added, which this post assumes. Part 3 covers a zero-downtime major upgrade.
The scenario: consolidation
This architecture is the one the documentation lists as “consolidating multiple databases into a single one, for example for analytical purposes”. The application developer’s version: three different applications (a shop, a CRM, a billing system), each with its own schema and its own server, all feeding a warehouse; and then the warehouse’s changes exported to something that is not Postgres.
The sources
Each application owns a schema named after it, on its own server. These are the three, trimmed to the tables that matter here (the demo has the full DDL and the rows):
-- on the shop server
create schema shop;
create table shop.customers
(
id int primary key,
account_id int not null,
name text not null,
email text, -- personal data, the warehouse will not get it
phone text,
country text,
tenant text not null
);
create table shop.orders
(
id serial primary key,
customer_id int not null references shop.customers(id),
amount numeric(10, 2) not null,
status text not null,
tenant text not null
);
create publication pub_shop for table shop.orders;
-- on the crm server
create schema crm;
create table crm.accounts (id int primary key, name text not null, tier text not null);
create table crm.contacts
(
id int primary key,
account_id int not null references crm.accounts(id),
name text not null,
email text
);
create publication pub_crm for tables in schema crm;
-- on the billing server
create schema billing;
create table billing.invoices
(
id serial primary key,
account_id int not null,
amount numeric(10, 2) not null,
status text not null
);
create table billing.payments
(
id serial primary key,
invoice_id int not null references billing.invoices(id),
amount numeric(10, 2) not null
);
create publication pub_billing for tables in schema billing;
The warehouse creates the same schemas and tables, by hand, because DDL is not replicated, and one subscription per application. Each subscription is owned by its own role, which will matter in a moment:
create role sub_shop login password 'x' in role pg_create_subscription;
grant create on database warehouse to sub_shop;
create schema shop authorization sub_shop;
create table shop.orders
(
id serial primary key,
customer_id int not null,
amount numeric(10, 2) not null,
status text not null,
tenant text not null
);
alter table shop.orders owner to sub_shop;
set role sub_shop;
create subscription sub_shop
connection 'host=shop dbname=shop user=repl password=repl'
publication pub_shop;
and the same for sub_crm and sub_billing. Once the initial copy is done,
the applications’ data sits side by side, and a query across applications is
plain SQL. Revenue per CRM account, from the CRM’s accounts and the billing
system’s invoices:
account | tier | invoiced | status
---------+--------+----------+--------
Acme | gold | 150.00 | paid
Globex | silver | 75.00 | open
Initech | bronze | 20.00 | open
A schema per application, and why it must start at the source
The warehouse keeps a schema per application, as you would want. But the way
it gets there is not what you might expect: a subscription cannot rename
anything. It looks for the publisher’s schema.table under the same name on
the subscriber. The subscriber’s code does a plain lookup with the names the
publisher sent (RangeVarGetRelid on the remote namespace and relation name,
in replication/logical/relation.c), and neither create subscription nor
create publication has an option to map one name to another.
That matters because the typical application does not have a schema of its
own: its tables are in public. I tried the natural thing. The shop has
public.orders, and the warehouse has shopapp.orders, where I want it.
The attempt:
create schema shopapp;
create table shopapp.orders (id int primary key, customer_id int not null, amount numeric(10,2) not null);
create subscription sub_rename
connection 'host=shop dbname=shop user=repl password=repl'
publication pub_rename;
ERROR: relation "public.orders" does not exist
The subscription is not even created. What works is to give the application a
schema of its own on the publisher, which is a lot less work than it sounds,
because a role’s search_path keeps the application’s unqualified SQL
resolving:
create schema shopapp;
alter table public.orders set schema shopapp;
create role app_shop login;
alter role app_shop set search_path = shopapp;
The publication follows the table, since it tracks it by identity and not by name, and the application does not notice, connecting with its own role and its own unqualified SQL:
pubname | schemaname | tablename
------------+------------+-----------
pub_rename | shopapp | orders
-- as app_shop: show search_path; select count(*) from orders;
search_path
-------------
shopapp
count
-------
2
Now the same create subscription as before finds shopapp.orders on both
sides, and the copy runs. If you have several applications in public on
several servers, this is the migration to do first, and it is cheap; if you
cannot touch the source, the alternatives are a database per source on the
warehouse server, or a component that renames as the changes flow, the kind
I come back to at the end of this architecture.
I also tried what happens when you do not do this. Two sources with a
public.customers and overlapping keys: the initial copy stops with
duplicate key value violates unique constraint "customers_pkey", and the table
stays in state d and retries every five seconds. Same name and a different
shape: logical replication target relation "public.contacts" is missing replicated column: "company".
Less data: filters and column lists
This warehouse is the EU warehouse. It must never hold the customers’ email
addresses and phone numbers, and it should only receive the eu tenant’s
rows. Since Postgres 15 the publisher does both, with a column list and a row filter.
The subscriber’s copy of shop.customers is created without the personal
columns, because a table only needs the columns that will be sent. On the
warehouse:
create table shop.customers
(
id int primary key,
account_id int not null,
name text not null,
country text,
tenant text not null
);
alter table shop.customers owner to sub_shop;
On the shop server, the publication is changed to list its tables with their
column lists and row filters. shop.orders was already published and gets a
filter, shop.customers is new:
alter publication pub_shop set table
shop.orders where (tenant = 'eu'),
shop.customers (id, account_id, name, country, tenant) where (tenant = 'eu');
select schemaname, tablename, attnames, rowfilter
from pg_publication_tables
where pubname = 'pub_shop'
order by tablename;
schemaname | tablename | attnames | rowfilter
------------+-----------+---------------------------------------+-----------------------
shop | customers | {id,account_id,name,country,tenant} | (tenant = 'eu'::text)
shop | orders | {id,customer_id,amount,status,tenant} | (tenant = 'eu'::text)
The first trap comes right after. The row filter uses tenant, which is not
part of the primary key, and the primary key is the replica identity. The
publication is accepted, and the next update on the publisher fails, whatever
column it changes:
update shop.orders set status = 'paid' where id = 4;
ERROR: cannot update table "orders"
DETAIL: Column used in the publication WHERE expression is not part of the replica identity.
Same fix as in the hub-and-workers demo: a unique index that contains the filter column, used as the replica identity.
create unique index orders_id_tenant on shop.orders (id, tenant);
alter table shop.orders replica identity using index orders_id_tenant;
create unique index customers_id_tenant on shop.customers (id, tenant);
alter table shop.customers replica identity using index customers_id_tenant;
Then the subscriber picks up the new table. refresh publication copies
only tables that are new to the subscription, and shop.customers is one, so
its initial copy honours the filter and the column list:
alter subscription sub_shop refresh publication;
id | account_id | name | country | tenant
----+------------+-------+---------+--------
1 | 1 | Alice | FR | eu
2 | 1 | Bob | FR | eu
4 | 3 | Dan | DE | eu
Carol, the US customer, is not there, and neither are the email and phone
columns. The second trap is in shop.orders, which was already being
replicated when the filter was added. A filter only applies to what is sent
from then on, so the US order that was copied before it existed is still on the
warehouse:
select id, tenant, status from shop.orders order by id;
id | tenant | status
----+--------+--------
1 | eu | paid
2 | eu | paid
3 | us | paid
4 | eu | paid
Row 3 stays until somebody deletes it by hand, and it will not follow later
changes either, since those are filtered out. After update shop.orders set status = 'shipped' where id = 3 on the shop, the warehouse’s row 3 still says
paid.
What the publisher sends
The column list is a promise that the personal columns never leave the
publisher, and I wanted to check it on the wire and not only on the subscriber.
A scratch logical slot with the pgoutput plugin, read with
pg_logical_slot_get_binary_changes(), returns the protocol messages the
subscriber would get. The function below lists the first byte of every message
(B begin, R relation, I insert, U update, D delete, C commit) and
whether any message contains an email or a phone number:
create function peek(slot text default 'peek', pub text default 'pub_shop')
returns table (messages text, pii_bytes boolean)
language sql as $$
select coalesce(string_agg(chr(get_byte(data, 0)), '' order by lsn), '(nothing)'),
coalesce(bool_or(encode(data, 'escape') ~ '(@example\.com|\+33|\+49|\+1 555)'), false)
from pg_logical_slot_get_binary_changes(slot, null, null,
'proto_version', '1',
'publication_names', pub)
$$;
select pg_create_logical_replication_slot('peek', 'pgoutput');
Then one statement at a time, followed by select * from peek();:
| Statement on the publisher | Messages | PII on the wire |
|---|---|---|
update shop.customers set phone = '…' where id = 1 (a column that is not published) | BRUC | no |
update shop.customers set name = 'Alicia' where id = 1 | BRUC | no |
update shop.customers set tenant = 'us' where id = 4 (the row leaves the filter) | BRDC | no |
update shop.customers set tenant = 'eu' where id = 3 (the row enters the filter) | BRIC | no |
insert into shop.orders … 'us' | nothing | no |
update shop.orders set status = 'shipped' where id = 3 (a row outside the filter) | nothing | no |
insert into shop.orders … 'eu' | BRIC | no |
Three things to read in that table. An update of a column that is not in the
list still sends an update message, only without the value: the subscriber gets
a no-op update. A row that leaves the filter arrives as a delete and a row
that enters it as an insert, so the warehouse follows the tenant changes. And
a change to a row outside the filter sends nothing.
To be sure the check can fail, a control: a publication on the same table without a column list, and the same phone update:
messages | pii_bytes
----------+-----------
BRUC | t
The warehouse ends up consistent with the filter, apart from the row that predates it:
id | account_id | name | country | tenant
----+------------+--------+---------+--------
1 | 1 | Alicia | FR | eu
2 | 1 | Bob | FR | eu
3 | 2 | Carol | US | eu
(Carol is eu now, because the demo moved her tenant. Dan, moved to us, was
deleted from the warehouse.)
Filters when the publication is a whole schema
The CRM’s publication is for tables in schema crm, which is convenient: a
new table in the schema is published without anybody touching the publication.
It comes with a limit. A column list is refused in a publication that has a
tables in schema element:
alter publication pub_crm add table crm.contacts (id, account_id, name);
ERROR: cannot use column list for relation "crm.contacts" in publication "pub_crm"
DETAIL: Column lists cannot be specified in publications containing FOR TABLES IN SCHEMA elements.
To hide the email of a contact you leave the schema form and list the tables, which gives up the “new tables follow automatically” behaviour:
alter publication pub_crm drop tables in schema crm;
alter publication pub_crm add table crm.accounts, crm.contacts (id, account_id, name);
tablename | attnames
-----------+----------------------
accounts | {id,name,tier}
contacts | {id,account_id,name}
And as with the row filter, it applies from now on. A contact inserted after the change arrives without its email, and the emails copied before stay where the initial copy put them:
insert into crm.contacts values (3, 3, 'Jane', '[email protected]');
id | name | email
----+------+---------------------
1 | Wile | [email protected]
2 | Hank | [email protected]
3 | Jane |
If the personal data must not be on the warehouse at all, the order matters: set the column list before the first copy of the table, or clean the subscriber up yourself.
Re-exporting as a change stream
The warehouse gets its rows through apply workers, and apply workers write WAL like any other session. So the warehouse tables are ordinary tables as far as logical decoding is concerned: a publication on them and a logical slot give a Debezium-like consumer the consolidated stream.
On the warehouse, one publication for the three schemas, and two slots created
with the client tools of the server: pgoutput, which is what Debezium uses,
and test_decoding, which is readable:
create publication cdc_pub for tables in schema shop, crm, billing;
pg_recvlogical -U postgres -d warehouse --slot cdc_pg --create-slot -P pgoutput
pg_recvlogical -U postgres -d warehouse --slot cdc_td --create-slot -P test_decoding
A slot belongs to one database, so a slot created in the postgres database
sees nothing of warehouse. Then some activity on the sources: one transaction
on the shop that inserts a customer and an order together, and later a write
made directly on the warehouse, to a payment. Reading the test_decoding slot
without consuming it:
select data
from pg_logical_slot_peek_changes('cdc_td', null, null,
'include-xids', '0',
'skip-empty-xacts', '1');
BEGIN
table shop.customers: INSERT: id[integer]:10 account_id[integer]:3 name[text]:'Ivy' country[text]:'FR' tenant[text]:'eu'
table shop.orders: INSERT: id[integer]:7 customer_id[integer]:10 amount[numeric]:42.00 status[text]:'new' tenant[text]:'eu'
COMMIT
...
BEGIN
table billing.payments: INSERT: id[integer]:900 invoice_id[integer]:3 amount[numeric]:1.00
COMMIT
Two things to read there. The applied rows do show up downstream, and the
customer’s email and phone are not in the stream, because the publication
that feeds the warehouse never sent them. And the source transaction that
touched two tables arrives as one downstream transaction.
The rows applied by a subscription carry a replication origin, named
pg_<subscription oid>, and the consumer chooses whether it wants them. The
same slot, read with pgoutput messages, with origin set to any and to
none (the letters are the message types: B begin, O origin, R relation,
I insert, U update, C commit):
select 'origin any' as option, string_agg(chr(get_byte(data, 0)), '' order by lsn) as messages
from pg_logical_slot_peek_binary_changes('cdc_pg', null, null,
'proto_version', '1', 'publication_names', 'cdc_pub', 'origin', 'any')
union all
select 'origin none', string_agg(chr(get_byte(data, 0)), '' order by lsn)
from pg_logical_slot_peek_binary_changes('cdc_pg', null, null,
'proto_version', '1', 'publication_names', 'cdc_pub', 'origin', 'none');
option | messages
-------------+-------------------------------
origin any | BORIRICBOICBOUCBORICBORICBRIC
origin none | BRIC
With origin = any the consumer gets every transaction, each applied one
preceded by an O message with the origin. With origin = none it gets only
BRIC, the payment written locally on the warehouse. test_decoding has the
same switch under another name: only-local. That is the trap: Postgres 16’s
origin filter, which is what protects you from loops, is also what hides
replicated data from a consumer that asks for it. A consumer of a consolidated
database must use any.
Large transactions stream, since Postgres 14. With logical_decoding_work_mem
at 64kB, a source transaction that inserts 20,000 rows reaches the downstream
slot while it is still open, with the option that asks for it
(stream-changes for
test_decoding). The rows are not visible on the warehouse yet, because the
source transaction has not committed:
downstream already streams while the source transaction is open: yes
rows tagged 'bulk_c' visible on the warehouse while the source transaction is open: 0
streamed changes delivered before the end of the transaction: more than 10000
rows tagged 'bulk_c' visible on the warehouse at the end: 20000
If the source rolls back instead, the downstream consumer has already received
thousands of changes, and the last line it sees is aborting streamed (sub)transaction. The consumer has to be able to throw them away. This demo
reads test_decoding’s output because it is plain text with a line per
change; pgoutput’s is a binary wire format, the same one earlier in this
post that only gave up its message types one byte at a time
(get_byte(data, 0)), so there is no text trace to show here for it. The
streaming protocol underneath is the same one Debezium already relies on in
production either way.
Keeping the CDC load off the primary
Since Postgres 16, a logical slot can live on a physical standby, which keeps the CDC consumer off the warehouse primary.
The standby is a base backup of the warehouse. -R writes the recovery
settings, and -C -S creates the physical slot on the primary:
pg_basebackup -h warehouse -U postgres -D $PGDATA -X stream -R -C -S standby1 -v
Creating the logical slot on it with pg_recvlogical needs four things, and I
collected them one error at a time:
The standby needs the primary’s sizing. With the default
max_worker_processesit does not start:FATAL: recovery aborted because of insufficient parameter settings,DETAIL: max_worker_processes = 8 is a lower setting than on the primary server, where its value was 24.It needs
wal_level = logicalitself, it does not inherit it from the primary, on every released version:ERROR: logical decoding requires "wal_level" >= "logical".Creating the slot blocks on an idle primary, until a running-transactions record reaches the standby. On the primary:
select pg_log_standby_snapshot();It needs
hot_standby_feedback = on. With it off, catalog vacuum on the primary removes rows the slot needs and the slot is invalidated:
slot_name | wal_status | conflicting | invalidation_reason
-----------+------------+-------------+---------------------
cdc_sb | lost | t | rows_removed
ERROR: can no longer access replication slot "cdc_sb"
DETAIL: This replication slot has been invalidated due to "rows_removed".
So the standby’s postgresql.conf carries three settings that the primary does
not force on it:
max_worker_processes = 24 # at least the primary's
wal_level = logical
hot_standby_feedback = on
Beta territory: one of those three may go away. Postgres 19 introduces
effective_wal_level, which raises the server’s actual WAL verbosity to
logical the moment any logical slot exists, without the wal_level GUC
ever changing. A separate primary and standby for this one check
(compose/sql/68-cdc-effective-wal-level.sh), both left at every default,
wal_level = replica included. A logical slot on the standby alone, before
the primary has any, is refused outright:
pg_recvlogical: error: could not send replication command "CREATE_REPLICATION_SLOT "probe_b" LOGICAL "test_decoding" ( SNAPSHOT 'nothing')": ERROR: logical decoding on standby requires "effective_wal_level" >= "logical" on the primary
HINT: Set "wal_level" >= "logical" or create at least one logical slot when "wal_level" = "replica".
One logical slot on the primary, nothing else reconfigured anywhere, is enough to raise it:
name | setting
---------------------+---------
effective_wal_level | logical
wal_level | replica
(2 rows)
the same two settings read on the standby itself, once its own slot exists too. And a row inserted on the primary decodes cleanly from the standby’s own slot, once it has replayed:
BEGIN
table public.t: INSERT: id[integer]:1
COMMIT
That removes the wal_level = logical line specifically, not the other
two: max_worker_processes sizing and hot_standby_feedback are a
different mechanism each, untouched by this. I also checked two other
19 candidates for this exact architecture, on purpose, and both turned out
not to help: EXCEPT on a publication refuses to combine with for tables in schema (ERROR: syntax error at or near "except"), so it does nothing
for the schema-wide publication limitation earlier in this post; and 18’s
new conflict counters stay at zero for the overlapping-keys failure from
“The sources”, because that failure happens during the initial table copy,
a different code path from the one the counters watch. Progress here is
real and it is also uneven — the next release does not improve everything
it touches, and that is worth knowing before you plan around a beta.
What breaks
DDL, again, and this time it is not a design decision you can avoid: the schema of the warehouse has to follow the schema of the sources by hand.
If the publisher adds a column that the subscriber does not have yet, on the shop:
alter table shop.orders add column note text;
insert into shop.orders (customer_id, amount, status, tenant, note) values (1, 5.00, 'new', 'eu', 'first with note');
insert into shop.orders (customer_id, amount, status, tenant) values (2, 6.00, 'new', 'eu');
the subscriber’s apply worker stops, and every later transaction queues behind the first one, including the second insert, which has nothing to do with the new column:
ERROR: logical replication target relation "shop.orders" is missing replicated column: "note"
subname | apply_is_failing
----------+------------------
sub_shop | t
orders_after_note_insert
--------------------------
0
slot_name | unconfirmed_wal
-----------+-----------------
sub_shop | t
The publisher keeps that WAL until the subscriber catches up, which is what the
last line shows. The fix is to add the column on the subscriber, and the worker
retries by itself at the next wal_retrieve_retry_interval (5 seconds by
default):
alter table shop.orders add column note text; -- on the warehouse
amount | note
--------+-----------------
5.00 | first with note
6.00 |
So the order for an added column is the subscriber first, then the publisher, and the same test with that order does not stop:
alter table shop.orders add column priority int; -- on the warehouse, then on the shop
Dropping a column goes the other way. The publisher goes first, and the subscriber keeps a column that is NULL for the new rows, until you drop it:
amount | priority
--------+----------
7.00 | 1
8.00 |
Why DDL isn’t replicated yet
Every tool that claims to replicate DDL does it one of two ways: fan the original command out directly, synchronously, to every node, the way Citus does from its coordinator; or capture it on the publisher and ship the text to be replayed on the subscriber, asynchronously, the way a logical replication setup like this one would need to. Core Postgres has had the building block for the second approach since 2013, and it is sitting there unused.
pg_logical_emit_message(transactional, prefix, message, flush) writes an
arbitrary blob straight to WAL, no table involved, and a subscriber that
asks for it (messages = true on the subscription) gets it over the wire
as an ordinary protocol message. It is exactly the shape a DDL statement’s
text needs: emit it from an Event Trigger on the publisher, decode it on
the subscriber, execute it there. I checked what core’s own apply worker
does with one of these messages today, and the answer is in the source,
verbatim:
case LOGICAL_REP_MSG_MESSAGE:
/*
* Logical replication does not use generic logical messages yet.
* Although, it could be used by other applications that use this
* output plugin.
*/
break;
Nothing requests the message in the first place — no code path in
CREATE SUBSCRIPTION ever sets messages = true — and if one arrived
anyway, this is what would happen to it: nothing.
pglogical’s replicate_ddl_command() does not use this path either. It
queues the command text in an ordinary table and relies on that table
itself being replicated; the apply side special-cases an insert into that
one table and executes its text instead of inserting the row — the same
trick Slony and Londiste used years earlier, not the WAL-message
mechanism. EDB’s Postgres Distributed, pglogical’s commercial successor,
replicates most DDL automatically, with no function call needed, but it’s
closed source and I couldn’t find a public description of the mechanism
underneath.
Citus doesn’t touch this problem in the same sense — there’s no
asynchronous gap to decode and replay across, since every worker is a
direct, synchronous connection from the coordinator — but “it just sends
the DDL” undersells what’s actually in src/backend/distributed/. A
citus_ProcessUtility hook intercepts every DDL statement before
Postgres runs it; a qualify pass rewrites every name in the parse tree
to be fully schema-qualified, so a worker’s search_path can’t change
what the command means; and a full deparser, one file per object type
under distributed/deparser/ (tables, types, functions, sequences,
roles, and more), turns the parse tree back into a command string for
each target rather than forwarding the client’s original text verbatim.
Table DDL gets a shortcut of its own: instead of deparsing a separate
statement per shard, Citus ships the one deparsed command to
worker_apply_shard_ddl_command(shardid, ddl_command) on each worker,
which reparses it locally and substitutes the shard-suffixed table name
into the parse tree before executing — the sharding rewrite happens on
the worker, not the coordinator. And “as part of the same transaction”
is not a figure of speech: a DDL statement that touches more than one
node goes through Postgres’s own two-phase commit, PREPARE TRANSACTION
on every worker from the coordinator’s pre-commit callback and COMMIT PREPARED from its post-commit one, with a maintenance-daemon recovery
path for whatever a crash leaves half-committed. It’s a deparser and a
rewriter too, just a synchronous one with no WAL in the loop.
Outside the Postgres-extension world the same table-and-trigger trick
keeps getting reinvented, not replaced. Xata’s open-source pgstream
captures DDL with an Event Trigger into its own pgstream.schema_log
table, for the same reason pglogical does: a downstream stream that
cannot afford to silently desync on a schema change. And a plain managed
Postgres service doesn’t raise the baseline either — Neon’s own docs on
logical replication are explicit that “the database schema and DDL
commands are not replicated,” and that adding or dropping a column still
has to be applied by hand on both sides, the same rule this post has
been living with throughout.
I have a personal stake in why the WAL-message path exists but nobody
has finished wiring it up for DDL. I wrote the original Event Trigger
patch, committed to Postgres 9.3 in 2012, and DDL replication was the
destination I had in mind for it: Event Triggers were the piece needed
to notice a DDL command, and I tried to ship the other half in the
same cycle — a deparser (ddl_rewrite.c) that would hand the trigger’s
function a normalized, replayable command string instead of just a
tag. Tom Lane’s response is the objection this whole section keeps
running into a decade later: “Why don’t you just pass the original
query string, instead of writing a mass of maintenance-requiring new
code to reproduce
it?”
That framing won. Event Triggers shipped with the command tag and the
original query string, not a deparser, and reconstructing a command
from its parse tree has stayed a separate, repeatedly attempted, never
finished project ever since. Álvaro Herrera picked that thread back up
in 2015 with “deparsing utility
commands”: a hookable CommandDeparse_hook and a ddl_deparse module
that takes the parsed command and its OID and turns them into a JSON blob,
with a matching function to render that blob back out as SQL text. That
code lived on the deparse branch of 2ndQuadrant’s BDR tree, and a small
piece of it is still in current Postgres today — test_ddl_deparse in
src/test/modules, explicitly documented as “not intended to do anything
useful on its own,” a unit-test fixture for the pg_ddl_command type, not
a production deparser. A much larger attempt to finish the job, a full
deparser plus the WAL-messaging infrastructure, went through more than
80 revisions on the mailing list
thread
between 2022 and 2024 before it was withdrawn. The problem was already
well understood before any of that
code was written. Asked in 2018 whether the community had a vision for
DDL replication, Peter Eisentraut’s answer still reads as the honest
summary of where things stand: “I think nobody has completely figured
this out yet. Whatever is in pglogical and bdr and similar external
projects are the best current compromises. But they have lots of
problems.”
Two other projects deparse SQL for entirely different reasons, worth
knowing about so you don’t reach for the wrong tool: pganalyze/pg_query
wraps Postgres’s own parser to normalise and fingerprint queries, which is
a read-side problem, not a replay-side one; and lacanoid/pgddl extracts
the DDL of objects that already exist, which is pg_dump’s job done as
callable SQL functions, not a live capture of a command as it runs. Neither
solves what this section is about.
Building pglogical’s trick with core only
pglogical’s mechanism is a table and a trigger. Nothing about that needs
the extension: it is buildable with a plain publication, a plain
subscription, and one fact about triggers that matters a lot here. An
apply worker runs with session_replication_role = replica, so an
ordinary trigger (ENABLE, the default) never fires for a replicated
change; only ENABLE REPLICA or ENABLE ALWAYS do. The demo below uses
ENABLE ALWAYS, which fires in both roles, so the same trigger also
replays a command written by hand, not just one that arrived through
replication. ENABLE REPLICA is the narrower, arguably safer choice for
this exact table: it fires only for the apply worker, so a plain local
insert into the queue table — arbitrary SQL, if this trigger is firing
at all — does nothing on its own.
On the shop server, a queue table and an Event Trigger that fills it from
ddl_command_end, capturing the literal text with current_query():
create table shop.ddl_log
(
id bigint generated always as identity primary key,
tag text not null,
command text not null,
logged_at timestamptz not null default now()
);
create function shop.log_ddl() returns event_trigger
language plpgsql as $f$
begin
insert into shop.ddl_log (tag, command) values (tg_tag, current_query());
end;
$f$;
create event trigger shop_log_ddl on ddl_command_end
when tag in ('CREATE TABLE', 'ALTER TABLE')
execute function shop.log_ddl();
On the warehouse, a matching table and the trigger that replays it,
created ENABLE ALWAYS so it fires in the apply worker, and owned by the
subscription’s own role for the same reason every other table in this
architecture has to be:
create table shop.ddl_log
(
id bigint primary key,
tag text not null,
command text not null,
logged_at timestamptz not null default now()
);
create function shop.apply_ddl() returns trigger
language plpgsql as $f$
begin
execute new.command;
return new;
end;
$f$;
create trigger apply_ddl before insert on shop.ddl_log
for each row execute function shop.apply_ddl();
alter table shop.ddl_log enable always trigger apply_ddl;
alter table shop.ddl_log owner to sub_shop;
Add the queue table to the publication, refresh the subscription, then run one real DDL statement on the shop, nothing special about it:
alter publication pub_shop add table shop.ddl_log; -- on the shop
alter subscription sub_shop refresh publication; -- on the warehouse
alter table shop.orders add column urgent boolean default false; -- on the shop
No manual alter table on the warehouse, this time. The column is there:
column_name | column_default
-------------+----------------
urgent | false
and the queue table doubles as an audit log of exactly what ran:
tag | command
-------------+-----------------------------------------------------------------
ALTER TABLE | alter table shop.orders add column urgent boolean default false
It is pglogical’s design, not a smaller version of it: current_query()
captures the statement as submitted, not a reconstruction, so it has the
same two edges pglogical has. A command sent as several statements in one
round trip is captured whole, not split; and it replays with whatever
search_path the apply worker’s session has, which is not necessarily the
one the original session had.
What pg_logical_emit_message actually sends
The other primitive from the previous section, proven rather than just cited. A second Event Trigger on the shop, emitting the same text as a logical message instead of a table row:
create function shop.emit_ddl_message() returns event_trigger
language plpgsql as $f$
begin
perform pg_logical_emit_message(true, 'ddl', current_query());
end;
$f$;
create event trigger shop_emit_ddl on ddl_command_end
when tag in ('CREATE TABLE', 'ALTER TABLE')
execute function shop.emit_ddl_message();
A scratch test_decoding slot needs no option to see it, and shows the
text verbatim:
select data from pg_logical_slot_peek_changes('peek_ddl', null, null)
where data like 'message:%';
message: transactional: 1 prefix: ddl, sz: 64 content:alter table shop.orders add column express boolean default false
pgoutput, what a real subscription uses, needs the option asked for
explicitly, confirmed by counting the message bytes on the wire with and
without it, on the same transaction:
select count(*) from pg_logical_slot_peek_binary_changes('peek_ddl_pg', null, null,
'proto_version', '1', 'publication_names', 'pub_shop')
where get_byte(data, 0) = ascii('M'); -- 0
select count(*) from pg_logical_slot_peek_binary_changes('peek_ddl_pg', null, null,
'proto_version', '1', 'publication_names', 'pub_shop',
'messages', 'true')
where get_byte(data, 0) = ascii('M'); -- 1
That’s as far as I built it. The message is really on the wire, with
messages = true, in the same transaction as the change it describes —
and, same as the queue table, only if that transaction commits: a
message emitted with transactional = true lives in the reorder
buffer with the rest of its transaction’s changes and is never decoded
at all if the transaction aborts. Turning the wire bytes into something
a subscriber acts on needs code that speaks the replication protocol
directly, the way pgcopydb already does for its own reasons — and
even there, the wire-level parser has no case for the Message frame
type yet; it falls through to the default branch and is logged as
unknown and dropped. Writing a general consumer, one that also has to
decode every row change pgoutput can send, is a real project, not a
demo, so this is where I stopped.
A narrower consumer is a smaller project, and worth naming precisely.
If the slot’s publication carries no tables at all, a client only ever
sees Begin, Message and Commit frames — none of pgoutput’s binary
tuple encoding, which is the genuinely hard part, ever shows up. That
alone makes “a small client that only replays DDL messages” a real
option, much smaller than a general consumer.
What it does not get for free is the one guarantee the queue table has by construction. A subscription applies its publisher’s transactions in the order they committed, one pipeline, one position in the WAL. Route the DDL message through a second, independent slot instead, and there are two pipelines reading the same WAL at their own pace, with no primitive in Postgres that makes one wait for the other — nothing stops the apply worker from inserting into a column the standalone DDL client hasn’t gotten around to adding yet. The queue table never leaves the pipeline: it’s an ordinary row, in the ordinary publication, applied by the ordinary apply worker, so the DDL inherits the ordering Postgres already promises instead of needing to rebuild it. A separate DDL-message slot would have to re-earn that guarantee some other way — a barrier table, an LSN watermark the apply worker waits on — which is coordination code on top of a smaller parser, not a net simplification.
Sequences don’t replicate either
DDL isn’t the only gap a subscription leaves for you to close by hand. Sequences don’t follow either, and this one isn’t a workaround-with-a- trigger kind of problem: the publisher’s sequence was at 40013, the warehouse’s copy at 1, and a local insert on the warehouse reuses an id that a replicated row already has:
ERROR: duplicate key value violates unique constraint "orders_pkey"
DETAIL: Key (id)=(1) already exists.
Postgres 19, still in beta, replicates sequences too — part 1 has the
details (FOR ALL SEQUENCES, on demand at REFRESH SEQUENCES, not
streamed) and the same caveat applies here: it is not yet something to
build this architecture’s correctness on.
That is where the tables-as-a-buffer approach shows its cost. Every change is written to the warehouse and then written again to be re-decoded. It works, and for many teams it is the right size. When it stops being the right size, you want a component that merges and splits change streams without writing tables, which is a story for another article.
Conclusion
This architecture has the same shape as the first one in this series, read
backwards: instead of one hub filtering reference data out to many workers,
it is many sources filtering their own data in to one warehouse. The same
releases did the work, walked the same way. Postgres 10 is still the base
pipe: a publication per source, a subscription per source, nothing else
required to get three applications’ data sitting side by side. Postgres
14’s streaming is what keeps a 20,000-row source transaction from making
the re-exported CDC stream wait for its commit. Postgres 15’s row filter
and column list are what keep the US order and the phone number off the EU
warehouse, in one alter publication, where an older stack needed a
trigger or a WHERE-clause on every extract job. Postgres 16 shows up
twice, for two different jobs: the origin option, the same mechanism
part 1 used to stop a two-way loop, here decides whether the re-exported
CDC stream carries the warehouse’s replicated data at all; and decoding
from a physical standby is what keeps that stream from loading the
warehouse’s own primary. None of it makes the one hard constraint go away
— a subscription cannot rename a table, so the naming has to be right at
the source — but each release made what is left of the problem smaller
and more declarative than the one before it.
The DDL rule from part 1 turns out to be the real constant across both architectures: the subscriber still has to move first for an added column and last for a dropped one, by hand, every time, in every one of these designs. Ten releases have not touched that, and nothing on the roadmap promises to.
pglogical deserves a direct mention here, not a footnote: it is a
production answer to both gaps this post worked around by hand. AWS
documents it as supported on every current RDS for PostgreSQL and
Aurora PostgreSQL release, replicating “sequences” by name alongside
tables — exactly the gap “Sequences don’t replicate either” hit above.
Azure Database for PostgreSQL ships it too, though its own docs are
upfront that DDL still isn’t automatic there either: you wire
pglogical.replicate_ddl_command() to an Event Trigger yourself, the
same shape this post built from first principles with core alone. EDB,
who maintains pglogical since acquiring 2ndQuadrant, still ships
compatibility releases — 2.4.6 added Postgres 17 and 18 — even though
the extension itself is in maintenance mode, its newer ideas going into
EDB Postgres Distributed instead. What this series tracks is narrower
on purpose: how far core Postgres alone gets, one release at a time,
without reaching for an extension at all. The DDL and sequence gaps are
exactly where that self-imposed boundary still shows, and on a cloud
that allows installing it, pglogical is very likely the faster way to
close both.
Part 3 of this series covers zero-downtime major upgrades with a way back. A fourth post, covering the architectures left out of this series in less detail — geo-replication, BDR-style multi-active setups, plain CDC and triggers — is also planned.
The demo is in the
compose/ directory
of this post’s source. Run it, break it, and tell me what I got wrong.
