Est.

Fan-In Replication From Multiple Databases to One Warehouse

Multiple independent databases require specialized handling to merge safely into one warehouse.

Staff Writer · · 11 min read
Cover illustration for “Fan-In Replication From Multiple Databases to One Warehouse”
Data Warehouses · October 7, 2026 · 11 min read · 2,408 words

Fan-in replication moves change data from several independent source databases into a single destination, and it introduces engineering problems that simply do not exist when one database feeds one warehouse. Three of them recur across every fan-in system: schema alignment, identity resolution, and cross-source ordering. Solving them requires deliberate choices at the capture layer, the merge layer, and the destination layer, not a bigger version of the same pipeline a single-source team would build.

Fan-in replication versus parallel single-source pipelines

The rise of microservices is the structural reason fan-in has become common. Each service owns its own database, so a query that once lived inside a single schema (a join between users and orders) now has to be reconstructed across engines that were never designed to talk to each other. Three scenarios make this unavoidable in practice. A company running users in PostgreSQL, orders in MySQL, and payments in MongoDB has to merge all three into one warehouse just to build a revenue dashboard. A retail business operating identical orders tables across PostgreSQL clusters in us-east, eu-west, and ap-southeast needs a single global_orders table in Snowflake to see its business as one business. A company migrating from Oracle to PostgreSQL has to run CDC from both systems simultaneously during the cutover window, feeding the same warehouse from two sources that will never fully agree with each other until the migration finishes.

Teams that treat fan-in as "add another connector" discover the cost later, usually in production. Schema conflicts between sources accumulate in the destination table. Duplicate records for the same real-world entity creep into customer tables. Events from different sources land in an order that doesn't reflect what actually happened. None of these failures throws an error. They sit quietly in the data until a dashboard reports a number that's wrong, or a model trains on a dataset that misrepresents its own history. The rest of this piece works through each of these three problems in turn, starting with how changes get captured from the source in the first place, since that is what makes any of them possible to fix.

Log-based CDC as the practical capture mechanism for fan-in architectures

Before a pipeline can merge streams, align schemas, or resolve identity, it has to capture changes from each source reliably. Log-based CDC reads a database's transaction log directly: the binlog in MySQL, the write-ahead log in PostgreSQL. This approach has minimal performance impact on the production system being read from, and it captures every change, including deletes, without requiring schema modifications or repeated polling of live tables. Query-based and trigger-based approaches to CDC cannot offer the same combination of completeness and low overhead. This is why log-based capture functions as the baseline assumption for any fan-in design.

That completeness matters more in fan-in systems than in single-source ones, because the risk compounds across sources. A capture mechanism that misses a delete on one source while correctly capturing deletes on another doesn't just produce one bad row. It produces an inconsistency between how the destination represents entities from source A versus source B, and that inconsistency grows with every additional source added to the pipeline.

Each engine exposes its log in its own way, and those differences shape what a fan-in pipeline has to handle downstream. PostgreSQL's logical replication decodes the write-ahead log into row-level change events, requires the wal_level parameter set to logical, works across major versions, and can replicate a defined subset of tables. MongoDB's Change Streams read from the oplog that the database already maintains for replica set replication, and they issue resume tokens that let a consumer restart from its last confirmed position without replaying events it already processed; this only works against a replica set or sharded cluster, since standalone MongoDB instances keep no oplog. DynamoDB sits apart from both: it does not guarantee exactly-once CDC delivery, so any fan-in pipeline that includes DynamoDB as a source has to build idempotency into its own downstream logic.

The three architectural patterns for merging CDC streams and their tradeoffs

Once changes are captured from each source, something has to decide how those separate streams become one destination table. Three patterns handle this merge, and the choice between them is a real architectural decision, not an implementation detail to settle later.

The first pattern fans each source into its own Kafka topic, with a stream processing layer, Flink, ksqlDB, or a custom consumer, reading across all of them, normalizing schemas, and writing a single unified stream to an output topic that the destination connector consumes. A regional consolidation pipeline under this pattern might run PostgreSQL in us-east into a topic called orders.us-east, PostgreSQL in eu-west into orders.eu-west, and PostgreSQL in ap-south into orders.ap-south, with Flink reading all three and writing into global_orders before it lands in Snowflake. This pattern gives a team the most control: schema normalization, deduplication, filtering, and enrichment can all happen before data ever reaches the destination. The cost is operational. Running and maintaining a stream processing cluster is real infrastructure, with its own failure modes and its own on-call burden. There's also a coordination risk specific to fan-in that occurs before the stream processing layer: if multiple CDC connectors and a nightly snapshot job all read from the same PostgreSQL instance independently, WAL retention can balloon or replication slots can build up faster than they drain. A platform that maintains a single connection to each source and fans data out downstream avoids that problem structurally, rather than requiring each consumer to manage its own connection to production.

The second pattern pushes transformation down to the connector level. Each source connector applies field renames, type conversions, and metadata injection as part of its own processing, then writes directly into the destination table without an intermediate stream processing cluster. This removes a significant piece of infrastructure from the pipeline, but it bounds how much normalization is possible by whatever the connector's transformation layer supports. Complex merge logic that goes beyond renaming and casting may still need to happen somewhere else.

The third pattern defers reconciliation to the destination itself. Each source writes into a staging or intermediate table, and a MERGE or UPSERT statement run against the warehouse reconciles all sources into the target table. This leverages the SQL engine the warehouse already provides for conflict resolution. Its tradeoff is latency: the destination is only as fresh as the last MERGE run, and schema conflicts have to be resolved before or during that MERGE. Each of these three patterns places the burden of schema alignment, identity resolution, and ordering in a different part of the system, and the next three sections work through each in turn.

Schema alignment across heterogeneous sources: where most fan-in pipelines first break

Schema mismatches are the most common reason fan-in pipelines fail, and they come from two separate causes: structural differences that exist at rest, and schema changes that happen over time without coordination across sources.

Structural mismatch is a near-certainty whenever sources are heterogeneous. Column names differ across engines even for identical concepts (user_id in one system, userId in another). Data types diverge in ways that look trivial until they aren't, such as PostgreSQL's timestamptz against MySQL's DATETIME. Document stores like MongoDB add a further wrinkle: with no enforced schema, field presence can vary from one document to the next within the same collection. None of this resolves itself. A transformation layer, whether that's Flink SQL in the stream processing pattern, field renaming at the connector, or a cast applied during a MERGE, has to normalize these differences before data reaches the destination. Skipping that step leaves the destination table structurally inconsistent, which corrupts every downstream query touching more than one source.

Fan-in also multiplies the damage any single schema change can do. In a single-source pipeline, an uncoordinated DDL change on that one database breaks that one pipeline, a contained and traceable failure. In a fan-in pipeline, an uncoordinated change on any one of several sources can corrupt the shared destination table for every source feeding it, since all of them write into the same merged structure. The practical fix is automatic schema change detection that adjusts the pipeline without manual intervention. A design that depends on an engineer noticing and reacting to every DDL event across every source does not scale once the source count moves past a handful.

Identity resolution: when the same entity lives in multiple source databases

Aligning schemas gets every source writing into the same structural shape, but it says nothing about whether two correctly structured rows from different sources describe the same real-world entity. Without resolving that question, the destination accumulates duplicate or split records that undermine any analysis joining across sources.

The clearest version of this problem is visible in building a 360-degree customer view. A customer might exist as a row in an authentication service backed by PostgreSQL, as a different row in a billing service backed by MySQL, and as yet another row in a CRM system. Each of these systems assigns its own primary key to that customer, with no shared reference between them. A naive merge produces three separate rows for one person. Resolving this requires a key strategy chosen before the pipeline is built. One approach matches on natural keys like email or phone number, with deduplication logic applied either in the transformation layer or at the MERGE step itself.

The regional consolidation scenario is a simpler case of the same underlying problem, because all three regional clusters share an identical schema and the same primary key space. The risk there isn't mismatched identity so much as key collision: order ID 1001 in us-east and order ID 1001 in eu-west are different orders that happen to share a number. The fix is a source_region metadata column injected at the connector level, tagging every row with where it came from before it ever reaches the shared table. This kind of metadata injection, a source tag, a region, a database name, is a standard connector-level transformation available under all three architectural patterns described earlier. What matters is that it gets applied consistently and at capture time. Managed CDC platforms that build transformations into the connector layer handle this metadata injection and field routing as part of normal processing, which removes the need to stand up a separate stream processing layer just to tag rows by source.

Cross-source ordering and the impossibility of a global event order

No global ordering exists across independent databases, and that is a property of distributed systems, not a limitation of any particular tool. A change committed in a US database and a change committed in an EU database at the same wall-clock moment have no inherent ordering relationship between them. Their system clocks may not agree within seconds of each other, so there is no reliable way to say which one "really" happened first. A pipeline that assumes it can reconstruct a single global sequence across sources will produce causally incorrect results, because the assumption itself doesn't hold.

Kafka topics don't fix this, and it helps to be precise about why. A single topic partition preserves order for events from one source, so all changes to a given key arrive from that source in the sequence they occurred. But there is no ordering guarantee across topics. Consuming from orders.us-east and orders.eu-west in parallel gives no assurance that the combined stream reflects the true chronological order of events across both regions, because no such combined order exists to reflect.

What's actually achievable is narrower, and it's enough. Per-key ordering within a single source holds reliably when Kafka partitions by primary key, guaranteeing that every change to a given row from that source arrives in sequence. Per-entity ordering across sources is a harder case: if the same entity key appears in multiple sources, ordering requires either a coordination point, a stream processor holding state per key, or an explicit acceptance that inter-source ordering is undefined and designing around that fact.

The practical response has three parts. First, inject a source_committed_at timestamp drawn from the source database's own log, not the time the event was ingested downstream, so that per-source ordering can be reconstructed later even if cross-source ordering can't. Second, use idempotent upserts keyed on a stable composite key combining entity ID and source tag, so that an event arriving out of order doesn't overwrite a more recent one with an older one. Third, for the migration cutover scenario specifically, accept that ordering between Oracle and PostgreSQL during the transition window is undefined by design. The right move there is to let the destination MERGE resolve conflicts on the composite key rather than attempting to engineer a sequencing guarantee that the underlying systems cannot support.

Exactly-once delivery in fan-in pipelines: what the guarantee covers

Exactly-once delivery is a guarantee about message processing, not about final data correctness, and the distinction matters most in fan-in pipelines because it's where that guarantee is most likely to be mistaken for something broader than it is. What exactly-once typically covers is that a given change event, once captured from a source's log, is processed by the pipeline and applied to the destination without being double-counted, even if the underlying transport retries delivery. That's a meaningful property. It is not the same as guaranteeing that the destination table correctly represents reality once sources, schemas, and entity identities multiply.

DynamoDB does not guarantee exactly-once CDC, so any fan-in pipeline that includes it as a source has to build idempotency downstream rather than assume the capture layer provides it. That idempotency, typically an upsert keyed on a stable composite key of entity ID and source tag, is what actually prevents duplicate processing from corrupting the destination, and it has to be designed deliberately.

Even where a pipeline achieves exactly-once processing end to end, that guarantee says nothing about whether the schema was normalized correctly, whether identity resolution correctly matched records across sources, or whether ordering assumptions held. A duplicate-free, exactly-once pipeline can still write a wrong answer into the destination if the entity resolution logic merged two customers incorrectly, or if a schema mismatch silently dropped a field. Exactly-once delivery is a necessary property of a well-built fan-in pipeline. It is not a substitute for the schema alignment, identity resolution, and ordering design that the rest of this piece has worked through, and treating it as one is how a technically correct pipeline still produces a warehouse full of wrong numbers.

Sources

  1. Simplifying data movement across multiple Clouds with richer CDC in Copy job in Fabric Data Factory Oracle source, Fabric Data Warehouse sink and SCD Type 2 (Preview)
Filed underData Warehouses

More in Data Warehouses