Fan-In Replication From Multiple Databases Into One Warehouse
Streaming from multiple databases into one warehouse creates schema conflicts and ordering problems.

Fan-in replication means streaming changes from several different production databases, each with its own engine and schema, into one shared analytical warehouse. The engineering problems this creates, schema conflicts, ordering guarantees, and pipeline coordination, have to get solved before a single row lands cleanly at the destination, and none of them are solved by simply running more pipelines. Fan-In Replication From Multiple Databases Into One Warehouse.
Fan-in replication and its architectural distinctness
Picture a company running Postgres for its order system, MySQL for billing, MongoDB for a product catalog, and DynamoDB for session state. Fan-in is the topology where many heterogeneous source databases stream their changes into a single analytical warehouse such as Snowflake, BigQuery, Databricks, or Redshift. This is a different problem than fan-out replication, where one source broadcasts to many replicas. In fan-out, the source dictates the schema and the replicas just receive it.
Teams end up in this position for an ordinary reason. Operational databases multiply because services and teams multiply, each choosing the engine that fits its own workload, and nobody centrally plans for how those systems will eventually need to talk to each other. That's the real-time data warehouse: a hub that collates changes from multiple sources without the delay batch processing would introduce.
The scale of this problem is not hypothetical. The global data replication software market was valued at $4.2 billion in 2025 and is projected to reach $9.1 billion by 2034, growing at an 11.3% compound annual rate, according to MarketIntelo, and more than 64% of organizations already replicate data across at least three environments, per IndustryResearch. Those numbers describe a landscape where fan-in is not a rare architecture reserved for unusually complex companies. It's becoming the default shape of the problem.
Fan-in is not N independent one-to-one pipelines running in parallel. It looks that way at first glance, source, connector, destination, repeated a few times. But the shared destination introduces ordering conflicts, schema collisions, and coordination requirements that simply don't exist when a pipeline runs from one source to one destination alone. This piece is an engineering breakdown of those problems.
How CDC captures changes at each source type
Every fan-in pipeline starts with change data capture, and there are three ways to do it, because the rest of this piece refers back to them. It's the gold standard for fan-in specifically because it's continuous and cheap at the source. Query-based CDC polls a timestamp column on a schedule, which is simple to set up but introduces lag tied directly to the poll interval and often requires adding a column that wasn't there before. Trigger-based CDC fires synchronously on every write, adding overhead to every transaction, and gets avoided in any high-write production workload for exactly that reason.
At fan-in scale, log-based capture shifts from a preference to a near-mandatory requirement. Polling overhead and trigger overhead both stack as sources multiply; log-based CDC doesn't, because it reads outside each source's write path regardless of how many other sources exist.
But naming the method doesn't resolve the heterogeneity, because each source type implements its log differently. Postgres exposes its WAL through logical replication, decoded via plugins like pgoutput, wal2json, or decoderbufs, and it requires the database configured with wal_level = logical, a Publication defined, and a Replication Slot held open. MySQL exposes a binlog. MongoDB requires a replica set or sharded cluster, standalone deployments don't support change streams at all, and consumption happens through an oplog-backed resume token. DynamoDB Streams holds changes for a fixed 24-hour retention window and enforces shard-level consumption constraints that shape how a consumer reads from it. SQL Server runs a dedicated capture job that reads the transaction log and writes into system CDC tables, with a separate cleanup job purging old records on its own schedule.
Before a single event reaches the shared warehouse, each source has already produced that event in a different format, under different guarantees, with a different failure mode waiting behind it. Fan-in has to reconcile all of that, at the same time, continuously.
The schema conflict problem when heterogeneous sources share a destination
In a one-to-one pipeline, a schema change at the source has exactly one path to follow: it propagates to the one destination table that mirrors it. Fan-in removes that simplicity, because multiple sources might describe the same logical entity in incompatible ways, or one source might change its schema entirely independent of the others sharing its destination.
Three categories of conflict occur repeatedly. The first is a naming collision: two sources both happen to have a table called "orders," but with different columns, different types, and different primary key strategies, and a naive fan-in setup merges them destructively, overwriting one source's structure with another's. The second is type mismatch: one source stores a field as VARCHAR(255), another treats the same logical field as TEXT or an ENUM, and the warehouse's coercion rules may not match what either source actually intended. The third is drift over time. A column gets added to one source's table but not the other's, and the warehouse schema has to evolve to accommodate that without breaking the pipeline feeding from the source that never changed.
Drift isn't a one-time migration problem. Production databases change constantly, columns added, renamed, dropped, retyped, on schedules no pipeline operator controls, and a fan-in pipeline that can't track that drift will produce wrong data without announcing that it's doing so.
The standard mitigation is namespace isolation: prefix or schema-qualify every destination table by its source, source_a__orders, source_b__orders, before any unification logic runs. That defers the harder question, how to actually join or merge these entities, to a transformation layer built to handle it deliberately, rather than letting an ingestion pipeline make that call implicitly and irreversibly.
Automatic schema evolution should be treated as table stakes rather than a premium feature. It means the pipeline detects a DDL change at the source and propagates it to the destination table without a human intervening and without dropping rows that are mid-flight when the change happens. Batch pipelines fail loudly when schemas mismatch; a job errors out, someone gets paged, the problem gets fixed before more damage happens. Streaming fan-in pipelines don't get that courtesy. A mismatch can silently truncate a column or coerce a type for hours before anyone notices, because the pipeline keeps running the entire time.
Ordering guarantees across independent change streams
Each source database keeps its own clock, its own transaction sequence, and its own delivery lag, and the warehouse receives all of it interleaved, with no global ordering guarantee holding the interleaving together.
That absence has real consequences for analytical correctness. Consider a join between an orders table replicated from Postgres and an inventory table replicated from MySQL. Nobody wrote a bug; the two streams just arrived out of step with each other, because nothing forced them into step in the first place.
Kafka partitioning is the standard tool for handling ordering within a single stream: partition by primary key, and order is preserved for a given entity inside that topic. But that guarantee stops at the topic boundary. Across topics, across sources, there's no ordering guarantee without deliberate coordination layered on top.
There's a second trap sitting under this one. Most message brokers default to at-least-once delivery. An event can be redelivered. In a single-source pipeline that's a minor inefficiency. In fan-in, a duplicated event from one source arriving after a conflicting event from a different source can silently overwrite correct state with stale state. Exactly-once semantics stop being a nice-to-have here and become a requirement.
Checkpointing is central to all of this. Each source connector tracks its own position independently, a WAL LSN for Postgres, a binlog offset for MySQL, a resume token for MongoDB, a shard iterator for DynamoDB Streams, and losing that checkpoint means re-reading from an earlier position and reapplying events that may now collide with events already applied from entirely different sources. The realistic target, given all of this, is per-entity ordering within each source rather than perfect global ordering. That's not achievable without a global coordinator sitting above every source, and nobody builds that in practice. The achievable target pairs per-entity ordering with idempotent upsert semantics at the destination, so that out-of-order arrivals across sources still converge on correct state eventually, even if they don't arrive in a tidy sequence. A join between orders (from Postgres) and inventory (from MySQL) in the warehouse may reflect orders that are newer than their corresponding inventory update, or vice versa, depending on which pipeline lagged, which matters for analytical correctness.
Pipeline coordination: fan-in is not N independent pipelines
The failure mode that appears first in practice is what amounts to a spaghetti architecture: separate point-to-point connections from each source into the warehouse, each managed on its own, forming a mesh nobody can observe as a whole. One pipeline can fail quietly, and nobody notices until the warehouse itself starts reflecting corrupt or incomplete state.
The fix is an event-driven intermediary sitting between the sources and the destination. A message broker, Kafka being the standard reference implementation for this layer, decouples each source connector from the destination writer entirely: sources publish independently of each other, and a single consumption model reads from all the resulting topics. Fan-in at that layer typically means Kafka topics organized by source table, with a stream processor, Flink being the standard reference for this kind of stateful work, consuming across multiple topics and applying exactly-once semantics before anything gets written to the warehouse.
Onboarding introduces its own coordination problem. Adding a new source to an already-running fan-in pipeline means applying that source's full initial snapshot to the warehouse without corrupting the state other sources have already written there. Incremental snapshots without table locks, a capability introduced in Debezium 1.6 and later, are the mitigation for this problem.
Failure recovery looks different at fan-in scale too. In a single-source pipeline, recovery is straightforward: replay from the last checkpoint and move on. In fan-in, a failed connector on one source can't be allowed to stall the destination writer for every other source feeding the same warehouse. The pipeline has to tolerate per-source lag and backpressure independently, or a hiccup in the least reliable source becomes an outage for the whole warehouse.
That requirement extends directly into monitoring. Fan-in needs per-source lag, not just an aggregate number that hides which source is actually behind, along with per-source error rates. It needs replication slot health for Postgres sources in particular, since Postgres retains WAL until a consumer confirms receipt, and a disconnected consumer will quietly fill the disk. It needs oplog window tracking for MongoDB sources, because if the oplog rotates past a stored resume token, a full re-snapshot becomes unavoidable. The DynamoDB Streams 24-hour retention window needs to be treated as a hard boundary. None of this is optional at scale. Per IndustryResearch's figures, 73% of enterprises prioritize replication investment specifically to maintain operational continuity during outages, and a fan-in pipeline running without solid observability is a single point of silent failure for an entire analytical estate.
Loading patterns at the warehouse under fan-in
Once events actually reach the warehouse, there are three ways to write them, and fan-in changes how each one behaves.
Append-only writes every CDC event as an immutable row, tagged with operation type and timestamp, from every source alike. It preserves full history without ever requiring a merge across streams, which makes it the pattern most naturally suited to fan-in, though querying current state out of it requires extra logic, deduplication, windowing, layered on top. Upsert, also called SCD Type 1, merges each event into a destination table by primary key, which is simpler to query but depends entirely on primary keys staying stable and non-conflicting across every source writing to that table, which is precisely where the same-table-name conflict described earlier turns from a modeling nuisance into a data-corruption risk. SCD Type 2 creates a new row for every change, with validity timestamps marking when each version was true, and fits dimension tables well, where knowing which source produced a value and when matters as much as the value itself.
Deletes deserve their own mention, because a delete event from one source must never be allowed to apply against a record that actually originated from a different source. Soft deletes, marking a row inactive rather than physically removing it, are the safer default across fan-in pipelines for exactly this reason.
The warehouse side of this has been moving. Databricks introduced AUTO CDC INTO in Lakeflow Declarative Pipelines by mid-2025, enabling SQL-based real-time change pipelines, and later reporting on the 2025 and 2026 release cycle describes meaningful latency and cost improvements for both SCD Type 1 and Type 2 handling. Microsoft's Fabric Data Factory Copy job, in preview and announced at FabCon/SQLCon 2026, added an Oracle CDC source alongside a Fabric Data Warehouse sink and SCD Type 2 support, though SCD Type 2 specifically isn't yet supported when Oracle is the CDC source, a limitation worth knowing before assuming the feature covers every combination.
Lakehouse targets built on Apache Iceberg carry their own hazard: streaming fan-in generates a large number of small files quickly, and merge-on-read versus copy-on-write isn't a default that happens to work at streaming scale, it's a configuration decision that has to get made deliberately, with compaction planned in from the start rather than bolted on after the small-file problem is already visible in query performance.
Not every table in a fan-in pipeline needs the same freshness. A fraud detection feed or an operational dashboard might genuinely need sub-minute updates, while a reporting table can tolerate a lag of minutes without anyone noticing. Applying uniform low latency across every table regardless of what it's used for is over-engineering, and it's a design smell.
What to evaluate in a fan-in replication tool
Every criterion here maps directly back to a failure mode already described, which is the only honest way to evaluate a tool: against the specific ways fan-in breaks, not against a generic feature checklist.
Source breadth and log-based coverage come first. Does the tool support log-based CDC for every source type actually in use, Postgres WAL, MySQL binlog, MongoDB change streams, DynamoDB Streams, SQL Server's CDC tables? A query-based connector anywhere in a fan-in estate is a latency and reliability liability that will eventually show up as the slowest, least trustworthy source in the pipeline.
Automatic schema evolution is the second criterion: does the tool detect and propagate DDL changes from any source without a human stepping in, and without dropping events that are in flight when the change lands? Third is exactly-once delivery, checkpointing and exactly-once semantics maintained per source independently, so a failure in one source's pipeline can't corrupt state already written from another. Fourth is per-source observability: lag, error rates, and health metrics exposed per source rather than folded into an aggregate number that hides exactly the source that's failing.
Fifth is operational overhead. Does the team have to build and run message broker clusters, connector configuration, and consumer code themselves, or does the tool absorb that infrastructure? Teams consistently underestimate the time it takes to make DIY CDC production-grade. Sixth is deployment speed. A tool that demands weeks of infrastructure work before the first row lands in the warehouse is, in a real sense, part of the problem it claims to solve.
The tool landscape spans a range of tradeoffs along these lines. Debezium, an open-source Kafka Connect plugin now on the 3.x release line active in 2026, supports Postgres, MySQL, MongoDB, SQL Server, Oracle, Db2, and others, with incremental snapshots introduced in 1.6 and exactly-once improvements added in 2.x releases, but running fan-in on top of it requires the team to own Kafka, connector configuration, and consumer code. AWS DMS is a managed option native to the AWS ecosystem, capable of heterogeneous migrations, though it carries known limitations around schema change handling and transformation capability, and it fits best when the destination is Redshift or S3 inside an already AWS-native stack, with limited portability outside it. Fivetran offers broad managed connector coverage suited to standard warehousing needs, though volume-based pricing can grow significant at high-throughput fan-in scale where latency isn't the binding constraint. Estuary is positioned, per a roundup from Guideflow, around low-latency CDC with a streaming-first architecture and usage-based pricing.
None of these tools eliminates the underlying engineering problems this piece has walked through. Schema conflict, ordering, and coordination don't go away because a vendor manages the infrastructure. What changes is who owns solving them, and how much of that solving happens before the first row ever lands.
Sources
- 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) | Microsoft Fabric Community
- Multi-Source CDC to a Single Destination: Merging Streams - Streamkap
- 7 best change data capture software for 2026 - Guideflow Blog


