Est.

Fan-In Replication From Multiple Databases to One Warehouse

Multiple databases feeding one warehouse demand their own failure modes and ordering guarantees.

Editor at Large · · 11 min read
Cover illustration for “Fan-In Replication From Multiple Databases to One Warehouse”
Data Warehouse · September 30, 2026 · 11 min read · 2,365 words

Streaming changes from several different databases, Postgres, MongoDB, DynamoDB, MySQL, into one warehouse is a different problem from ordinary replication, with its own failure modes. It is a different problem with its own failure modes. A single-source pipeline rests on three quiet assumptions: failures stay local, one schema governs the destination table, and one transaction log defines the order events happened in. Fan-in knocks out all three assumptions at once, not one at a time.

That matters because of what fan-in gets used for. Operational dashboards, AI agents making live decisions, fraud detection, A/B testing: every one of these depends on data that is both fresh and internally consistent, and every one of them produces a wrong answer, not just a slow one, when the pipeline feeding it gets confused about ordering or schema. The stakes rise as the pattern spreads. Most enterprises now replicate data across several environments at once, spanning cloud, on-premise, and edge infrastructure. Multiply a single-source pipeline's list of failure modes across that many sources and environments, and the naive assumption that fan-in is just CDC with more connectors attached stops holding up. It requires its own mental model, built around what can and cannot be guaranteed when independent systems feed a shared destination.

Log-based CDC across heterogeneous sources

The three databases that show up most often in fan-in setups, Postgres, MongoDB, and DynamoDB, each expose change events through a mechanism built for that database alone. None of them were designed with each other in mind, and that mismatch is the root of almost everything that follows.

Postgres logical replication reads the write-ahead log (WAL) and turns it into a stream of row-level inserts, updates, and deletes, exposing logical changes instead of the raw physical page writes underneath. A connector reads that stream as a logical replication client, pulling from a replication slot that keeps precise track of how far it has read. The database has kept improving this path: Postgres 15 added row and column filtering on publications, Postgres 16 allowed decoding directly from a standby, and Postgres 18, released September 2025, made parallel streaming the default and added conflict logging. The risk sits in that replication slot. It holds onto WAL segments until a connector actually reads them, so a connector that stays down for an extended stretch causes WAL to pile up on disk, and in the worst case fills the disk and takes the source database down with it.

MongoDB works differently. It exposes Change Streams, which cover insert, update, replace, delete, drop, rename, dropDatabase, and invalidate events, and from MongoDB 6.0 onward, additional DDL events including create, createIndexes, and dropIndexes. Running this in production means paying attention to replica set configuration and oplog window sizing, because resume tokens have a shelf life: if the oplog rolls over before a paused connector comes back, that resume token stops working and the connector has no way to pick up where it left off.

DynamoDB takes a third approach entirely, using DynamoDB Streams with configurable stream view types and shard iterators to walk through changes. B|Its retention window is the tightest of the three, and a connector down for longer than the retention window loses those changes permanently. Shard splits need explicit handling in the connector logic, and IAM permissions control who can even read the stream in the first place.

Log-based CDC beats query-based or trigger-based approaches for warehouse feeds on every count that matters: it barely touches the source database, it captures deletes along with inserts and updates, and it doesn't require any schema changes on the source to work. That's why it has become the default mechanism behind fan-in pipelines. But three different logs, three different retention models, and three different recovery semantics, running under one shared destination, is where the real engineering work begins.

The ordering problem: why global event sequencing across sources is unsolvable in the simple case

Each source's log gives a clean, total order within itself. Postgres WAL LSNs, MongoDB resume tokens, and DynamoDB shard iterators all guarantee that events from a single source arrive at the warehouse in the sequence they actually happened. None of them say anything about how events from different sources relate to each other in time, and that gap cannot be closed by adding more infrastructure.

Network delay can compound ordering problems, because an event that occurred first at the source can arrive at the warehouse second if its path through the pipeline was slower. Routing everything through a shared broker, Kafka topics per source, does not solve it either. Merging separately-ordered streams into one globally ordered stream requires forcing every event through a single serialization point, which turns the broker into a bottleneck and defeats the reason for partitioning by source in the first place.

The upsert pattern, a MERGE keyed on primary key, hides this problem rather than solving it. If a delete from one source and an insert from another arrive out of sequence for the same key, the merge produces whatever state happened to land last, regardless of which event actually happened last in reality. That is not a rare edge case in a fan-in system with several active sources.

Warehouse tables fed by fan-in should be built around what the system can actually promise: order within a source, not order across sources. That points toward append-only tables carrying operation_type and source_transaction_id metadata, so every row records its origin and its position in that source's own sequence, and the full history stays queryable and auditable rather than silently overwritten. On the pipeline side, that means partitioning Kafka topics by source and primary key, monitoring lag separately for each source, and being explicit with downstream consumers about the guarantee they're actually getting: per-source total order, never a global one.

Namespace collisions and table naming conflicts when multiple sources feed one destination

Microservice architectures produce a lot of tables named users, orders, and events, and that repetition turns into a real problem the moment several of those services feed the same warehouse. Writing them all in without a namespace plan produces either a silent overwrite or a merge that quietly corrupts both.

The collision is not always obvious: two sources may share a table name but have different schemas, different primary key definitions, or different business meanings for the same column name. Because that corruption produces no error at the time it happens, a query may return something wrong only well after the fact.

Prefixing by source (service_a.users, service_b.users) keeps each source isolated but forces every consumer to know in advance which source they need, and makes cross-source joins clumsy. Separate schemas per source give clean isolation at the warehouse layer, and both Snowflake and Redshift support this natively, while BigQuery achieves the same effect through datasets; the tradeoff is schema management overhead that grows with every new source added to the pipeline. A single canonical table with a source_id column enables cross-source analysis directly, but only holds up if every contributing source keeps a compatible schema, which is precisely the assumption schema drift tends to break.

Primary key conflicts make this worse in a way that's easy to miss until it happens. Two sources can each have a row with id = 42 referring to two completely unrelated entities, and an upsert merge run without source scoping will overwrite one with the other and never raise an error. Debezium's topic.prefix field gives a source-scoped namespace at the broker level, but that scoping has to be carried all the way through to the warehouse table or schema name. It is a small step, and one easy to skip in custom consumer code, and that omission is why so many fan-in pipelines end up with this failure mode.

Schema evolution across sources: how a column change in one database breaks the warehouse for everyone

A schema change on any single source in a fan-in system propagates straight into a shared destination, and because every source evolves on its own release schedule with no coordination between them, the warehouse has to absorb incompatible changes arriving from multiple directions at once.

In an ordinary single-source setup, a column that gets added or retyped upstream can already break a pipeline, sometimes without any visible error at all. Fan-in raises the stakes on that same failure: the blast radius now covers every consumer reading the shared warehouse table, not just the one pipeline tied to the source that changed. A schema break in a single-source pipeline is an inconvenience. The same break in fan-in is a multi-team incident.

The two standard warehouse strategies handle this pressure differently. SCD Type 1, upserting in place, is fast and simple to run day to day, but a column rename or type change on one source can produce a merge that's incompatible with the existing warehouse schema, silently truncating, coercing, or dropping data without any alert. SCD Type 2, appending with versioning, preserves history and absorbs additive changes without much trouble, but a genuinely breaking type change on one source still forces a schema migration, and that migration can invalidate historical rows contributed by every other source sharing the table.

Two mitigations do real work here: maintaining a materialized view at the source that exposes only the fields downstream consumers actually need decouples the CDC pipeline from the source's internal data model, so an internal refactor on that source never has to touch the warehouse at all. Schema Registry compatibility rules, BACKWARD, FORWARD, FULL, add a contractual check at the broker: a connector attempting to publish a breaking change fails validation before that change ever reaches the warehouse. The usual objection is that a Schema Registry adds operational overhead to maintain. In a fan-in system, the cost of a silent schema break hitting a table shared across multiple consumers runs higher than the cost of keeping compatibility rules enforced, so the tradeoff favors enforcement. Postgres offers a complementary, source-side option: row and column filtering on publications, introduced in Postgres 15, limits which columns even enter the replication stream, shrinking the surface area available for a breaking change to reach the pipeline in the first place.

Failure isolation: keeping one broken source from taking down the whole pipeline

Diagram: Three Sources, Three Retention Clocks. Visualizes: Visualize the three CDC source mechanisms and their retention/recovery constraints side by side to show how each has a different failure cliff.

A connector failure or a WAL buildup on one source threatens every other source feeding the same warehouse, unless the architecture draws hard lines between failure domains from the start.

Postgres illustrates the risk clearly. If a connector tied to one Postgres source goes down and isn't restarted within an acceptable window, its replication slot keeps accumulating WAL with no upper bound, and in extreme cases that fills the disk and takes the source database itself offline, even though the MongoDB and DynamoDB pipelines running alongside it were never at fault. DynamoDB's failure mode looks different but is arguably less forgiving: its 24-hour retention ceiling means a connector down for more than a day loses those changes permanently, with no recovery path at all, independent of how the rest of the system is performing.

Containing that risk means giving each source its own connector instance, its own Kafka topic namespace (using a topic-prefix convention), and its own consumer group, so a failure on one source cannot block progress on the others. The preferred pattern for running this at scale is federated control paired with centralized observability: each source manages its own replication state independently, while one monitoring layer surfaces lag, error rates, and health for every source in a single view. Skip that centralization and teams end up with blind spots, discovering a stalled source only after it has already caused damage.

For Postgres specifically, the operational safeguards are concrete: cap WAL retention with a max size setting, adjust transaction log retention ahead of any planned downtime, and set alerts that fire once a connector has been down longer than an acceptable threshold. None of this replaces good monitoring elsewhere, but it closes off the failure mode that does the most damage the fastest.

Isolation only solves half the problem, though. Once a connector restarts, it replays events from its last committed offset, and the warehouse has to handle that replay without creating duplicate rows. That requires either exactly-once delivery guarantees or deduplication logic built into the destination, and that requirement is where the next challenge starts.

Exactly-once delivery in fan-in

Exactly-once delivery in a single-source pipeline is already a demanding property to hold onto. In fan-in, it has to hold independently for every source at once, and the warehouse has to manage several concurrent exactly-once guarantees landing on the same destination table simultaneously.

Achieving exactly-once delivery requires two separate hops: source to broker, and broker to destination, each needing its own guarantee. Getting exactly-once from source to broker depends on transactional guarantees at the broker layer, Kafka's transactional producer being the standard mechanism. Getting exactly-once from broker to destination depends on the destination connector tracking its own internal state to deduplicate, or on the warehouse performing that deduplication itself during the merge. Neither hop implies the other; both have to be handled, for every source, at the same time.

The tooling has moved to make this more tractable. As of 2025, Debezium 2.x added improved exactly-once semantics support along with incremental snapshots that don't require table locks, and Kafka 4.0's KRaft mode removes ZooKeeper from the coordination path entirely, cutting out one more place where things could go wrong. Apache Flink CDC Connectors, donated to Apache as Flink CDC 3.x in 2024, take a different approach: they provide exactly-once processing with stateful operations built in natively, without needing Kafka as a middle layer at all, which reduces the number of components exactly-once guarantees have to be coordinated across.

Fan-in adds one more wrinkle to deduplication that single-source pipelines never have to face. A duplicate event from one source and a genuine new event from a different source can look identical once they reach the warehouse, if the namespace strategy covered earlier doesn't carry through into the deduplication key itself. Source scoping has to run end to end, from the connector all the way to the merge logic; skipping it at the broker, where it's easiest to set up and easiest to forget, causes problems everywhere downstream.

Sources

  1. CDC for Real-Time Data Warehousing | Conduktor
  2. Database Replication Patterns: Active-Active, CDC, and Beyond - Streamkap
  3. Change Data Capture in 2026: What It is and How it Works?
  4. What is Change Data Capture? CDC Fundamentals | Conduktor
  5. CDC (change data capture)—Approaches, architectures, and best practices
Filed underData Warehouse

More in Data Warehouse