Dev & EngARTICLE

How to consolidate multiple databases into a single Postgres with logical replication

In the second part of his series on logical replication, Dimitri Fontaine shows how to merge schemas from different applications into a single Postgres warehouse, and still republish the changes as a CDC stream.

The scenario: three applications, one warehouse

French consultant Dimitri Fontaine published the second part of a series on logical replication on Planet PostgreSQL, a series that already spans ten years and ten versions' worth of accumulated features. The first part covered scaling writes with a hub-and-workers architecture; the third promises to cover major version upgrades without downtime. This one covers the case that Postgres's own documentation cites as an example: consolidating multiple databases into a single one, for analytical purposes.

The practical example brings together three applications typical of any company: a store (schema shop), a CRM (schema crm), and a billing system (schema billing), each on its own server with its own schema. Each server creates a publication for its relevant tables, and the warehouse creates a subscription per application, with a dedicated role per subscription. After the initial copy, the data from the three applications sits side by side in the warehouse, and cross-referencing CRM revenue with billing becomes a plain SQL query, with no external ETL.

Diagram shows three application servers, each with its own schema, publishing via logical replication to a central warehouse database, which in turn republishes its own changes through a logical slot to a Debezium-style consumer
Diagram shows three application servers, each with its own schema, publishing via logical replication to a central warehouse database, which in turn republishes its own changes through a logical slot to a Debezium-style consumer. Reprodução: postgr.es.

The schema can't be renamed: the fix starts at the source

The most important point in the article, and what Fontaine calls a trap, is that a subscription has no way to map names: it looks, at the source, for exactly the schema.table that exists at the destination. Anyone who tries to keep the table in public on the application side and shopapp.orders on the warehouse gets a direct error: ERROR: relation "public.orders" does not exist, and the subscription never even gets created.

The solution he demonstrates is moving the table to its own schema on the publisher, not on the subscriber:

sql
alter table public.orders set schema shopapp;

create role app_shop login;
alter role app_shop set search_path = shopapp;

With the role's search_path adjusted, the application's unqualified SQL keeps resolving normally, and the publication keeps following the table because it tracks object identity, not name. Anyone who skips this step and tries to bring up two sources with identically named public.customers tables gets duplicate key value violates unique constraint during the initial copy; different schemas with the same table and diverging columns fail with missing replicated column. In both cases, the problem isn't the tool: it's modeling that should already have been solved before any replication came into play.

Less data on the wire: column lists and row filters

Since Postgres 15, a publication can restrict both columns and rows. In the example, the European warehouse can't receive customer email and phone numbers, nor rows for customers outside the European Union:

sql
alter publication pub_shop set table
 shop.orders where (tenant = 'eu'),
 shop.customers (id, account_id, name, country, tenant) where (tenant = 'eu');

The destination table in the warehouse is created without the sensitive columns, because it only receives what will be sent. Fontaine verified this directly at the protocol level, not just by trusting the documentation: he wrote a function that reads pg_logical_slot_get_binary_changes() with the pgoutput plugin and inspects the bytes of each message looking for email and phone patterns.

Action on the publisherMessages in the slotPII on the wire
Update to an unpublished columnBRUCno
Update to a published column, no filterBRUCyes (control)
Row that enters the filterBRIC (becomes insert)no
Row that leaves the filterBRDC (becomes delete)no

The behavior of a row entering or leaving the filter is the detail that matters most to anyone operating this in production: Postgres converts the tenant change into a logical insert or delete, so that the warehouse tracks the row entering and leaving the slice.

The replica identity catch

Filtering by a column that isn't part of the primary key breaks replication on the very next change. Since the filter uses tenant, and the default replica identity is the primary key itself (without tenant), any update on the publisher fails with Column used in the publication WHERE expression is not part of the replica identity. The fix is the same as in the series' first article: create a unique index that includes the filter column and declare it as the replica identity.

sql
create unique index orders_id_tenant on shop.orders (id, tenant);
alter table shop.orders replica identity using index orders_id_tenant;

For anyone who has administered databases long enough, this requirement shouldn't come as a surprise: a logical filter on a column is, in practice, a query predicate, and every query predicate needs an index to back it. It's the same reasoning behind execution plans, applied to replication.

What was already there doesn't move on its own

Another point the text makes explicit: alter subscription ... refresh publication only copies new tables into the subscription. A filter added after a table was already replicating doesn't apply retroactively: a row that was copied before the filter existed stays in the warehouse and stops receiving updates, because subsequent changes are discarded at the source. Fontaine shows this with an order placed after the filter was added: the old status (paid) stays in the warehouse even after the source changes to shipped. Anyone who depends on consistent data needs to clean up these rows manually or set up the filter before the initial copy.

Limitation when the publication covers an entire schema

A publication of the for tables in schema type is convenient because new tables are picked up automatically, without changing the definition. The price is that it doesn't accept a column list: trying to restrict columns on crm.contacts within that publication returns an error, and the only way out is to abandon the schema-based form and list the tables manually, losing the automatic behavior for future tables.

Re-exporting the warehouse as a CDC stream

The part most relevant to anyone building data pipelines comes next: the application workers that write to the warehouse generate WAL like any other session. This means the warehouse can have its own publication and a logical slot read by a Debezium-style consumer, mixing the data that came from the three sources with what's written directly there.

The detail that catches anyone who doesn't read the fine print is the origin option, available since Postgres 16. It exists to avoid replication loops, but it also filters out replicated data by default: with origin = none, the consumer only sees what was written locally in the warehouse; with origin = any, it sees everything, including what arrived via the subscriptions. A consumer of a consolidated database needs to explicitly request any, or it will think the pipeline is broken when it's actually filtering by design.

Large transactions arrive before the commit

Since Postgres 14, large transactions stream to the logical consumer even before they're committed at the source, controlled by logical_decoding_work_mem. In the test described, a twenty-thousand-row transaction arrives at the slot in chunks while still open at the source, even though those rows aren't yet visible in the warehouse. If the original transaction is rolled back, the consumer receives an aborting streamed (sub)transaction message and has to discard everything it had already processed. It's a behavior that any consumer based on this protocol, Debezium included, already has to handle, but one that usually only shows up at production volumes, not in small tests.

Taking the CDC load off the primary

Also since version 16, a logical slot can live on a physical replica, not just on the primary. The standby reports its needs to the primary via hot_standby_feedback, so the primary preserves the catalog rows the slot needs to decode. Creation is done with pg_basebackup, combining -R (writes the recovery configuration) with -C -S (creates the physical slot on the primary):

pg_basebackup -h warehouse -U postgres -D $PGDATA -X stream -R -C -S standby1

This removes the cost of CDC logical decoding from the primary, freeing it up for the transactional workloads that actually need it.

When it's worth it, and when it isn't

In short: this design works well when the number of sources is small, known, and stable, and when the team wants a consolidated database queryable with plain SQL, without depending on an external streaming layer. It does, however, require manual intervention at points that a traditional ETL pipeline would solve differently: DDL doesn't replicate, so schema evolution needs explicit coordination between source and destination; renaming a table or moving it to a schema is a change that starts at the application, not at the warehouse; and filters, column lists, and PII exclusions only take effect from the moment they're configured, never retroactively.

For anyone evaluating this architecture, the underlying lesson in Fontaine's article is broader than any configuration flag: before thinking about volume, network round-trips, or CDC throughput, the modeling needs to be solved at the source. A poorly designed schema in public, a primary key that doesn't cover the business filter, or a missing replica identity aren't implementation details, they're the root cause behind every replica that stops working.

Translated from the Brazilian Portuguese original · Read the original