PostgreSQL · Architecture

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.

Three application servers, one schema each, subscribe into one warehouse database that keeps a schema per application. The warehouse then publishes its own change stream, read through a logical slot by a Debezium-like consumer. Three application servers, one schema each, subscribe into one warehouse database that keeps a schema per application. The warehouse then publishes its own change stream, read through a logical slot by a Debezium-like consumer.

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 publication sits between the two tables. The column list drops email and phone, the row filter drops the us customer. Neither the columns nor the row are ever sent to the warehouse. The publication sits between the two tables. The column list drops email and phone, the row filter drops the us customer. Neither the columns nor the row are ever sent to the warehouse.

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 publisherMessagesPII on the wire
update shop.customers set phone = '…' where id = 1 (a column that is not published)BRUCno
update shop.customers set name = 'Alicia' where id = 1BRUCno
update shop.customers set tenant = 'us' where id = 4 (the row leaves the filter)BRDCno
update shop.customers set tenant = 'eu' where id = 3 (the row enters the filter)BRICno
insert into shop.orders … 'us'nothingno
update shop.orders set status = 'shipped' where id = 3 (a row outside the filter)nothingno
insert into shop.orders … 'eu'BRICno

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.

The apply workers write the changes of the three sources into the warehouse. Logical decoding reads them from the WAL like any other change, together with the writes made directly on the warehouse. The origin option of the consumer decides which of the two it gets. The apply workers write the changes of the three sources into the warehouse. Logical decoding reads them from the WAL like any other change, together with the writes made directly on the warehouse. The origin option of the consumer decides which of the two it gets.

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 logical slot lives on the physical standby and decodes the WAL the standby replays. The standby reports its needs to the primary with hot_standby_feedback, so that the primary keeps the catalog rows the slot needs. The logical slot lives on the physical standby and decodes the WAL the standby replays. The standby reports its needs to the primary with hot_standby_feedback, so that the primary keeps the catalog rows the slot needs.

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_processes it 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 = logical itself, 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.

DDL is not replicated. For an added column the subscriber gets it first, otherwise its apply worker stops. For a dropped column the publisher goes first, and the subscriber keeps a column that stays NULL. DDL is not replicated. For an added column the subscriber gets it first, otherwise its apply worker stops. For a dropped column the publisher goes first, and the subscriber keeps a column that stays NULL.

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.

Event Triggers in 2012 made DDL visible; ddl_deparse in 2015 turned it into text, on an external tree; a full deparser and WAL-messaging patch went through 80+ revisions between 2022 and 2024 and was withdrawn. Today pglogical, PGD, Citus and Xata's pgstream each replicate DDL their own way, and core's own WAL-message primitive still has no consumer. Event Triggers in 2012 made DDL visible; ddl_deparse in 2015 turned it into text, on an external tree; a full deparser and WAL-messaging patch went through 80+ revisions between 2022 and 2024 and was withdrawn. Today pglogical, PGD, Citus and Xata's pgstream each replicate DDL their own way, and core's own WAL-message primitive still has no consumer.

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.

The DDL statement fires an Event Trigger that inserts into a queue table on shop. The row replicates like any other row. On warehouse, the apply worker's insert fires an ENABLE ALWAYS trigger that executes the command text — the same table-and-trigger shape Xata's pgstream uses too. The DDL statement fires an Event Trigger that inserts into a queue table on shop. The row replicates like any other row. On warehouse, the apply worker's insert fires an ENABLE ALWAYS trigger that executes the command text — the same table-and-trigger shape Xata's pgstream uses too.

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.

Dimitri Fontaine
Dimitri Fontaine
PostgreSQL Major Contributor & author of The Art of PostgreSQL