Real-Time Data Ingestion: Architecture, Methods & How It Works
A food delivery app has an ops dashboard that flags any order stuck in preparing for more than 15 minutes. If the dashboard reads from a warehouse that loads overnight, it flags those orders the next morning, long after the customers ordered somewhere else.
Real-time ingestion closes the gap between a row changing in Postgres and that change being queryable downstream. This post is for data engineers and platform teams weighing how to close it. It covers the methods, the parts of a working ingestion pipeline, and what breaks in production.
Artie is a fully managed CDC platform that captures row-level database changes and maintains destination tables in Snowflake, BigQuery, Redshift, Databricks, and other supported destinations.
Key Takeaways
- Real-time ingestion captures and delivers changes continuously instead of on a schedule.
- Micro-batching shortens the collection window but keeps it.
- Log-based change data capture (CDC) is usually the preferred method for production databases when low-latency, complete row-change capture is required.
- Event and API ingestion cover event-native data.
- A working architecture has up to six parts, and Kafka is common but not required.
- If nothing acts on the data faster than the batch interval, batch is cheaper and simpler.
- The hard parts show up in production: WAL growth, schema changes, backfills, and duplicates.
What Real-Time Data Ingestion Actually Means
Ingestion moves data from where it's produced into a system where it can be queried. Three models differ in when that move happens.
| Model | When data moves | What sets freshness |
|---|---|---|
| Batch | On a schedule, in bulk | The schedule, such as hourly or nightly |
| Micro-batch | In small windows of seconds to minutes | Window length plus load time |
| Streaming | Continuously, with no schedule | Capture, transport, and write time |
Real-time data ingestion means changes reach the destination within seconds to a minute of being committed, without waiting for a schedule. Streaming data ingestion sits at the continuous end of that range.
The difference between micro-batch and streaming is where the waiting happens. A micro-batch system collects records until a trigger fires, then processes them as a group. Spark Structured Streaming works this way by default. A streaming system reads the source continuously and delivers changes as they arrive, with no schedule gating capture.
Streaming doesn't mean every record is written the instant it commits. Warehouses handle bulk writes far better than row-by-row inserts, so many streaming tools capture continuously and apply changes in short batches lasting seconds, not tied to a schedule.
When a courier marks an order picked_up, batch reports it tomorrow, a 15-minute micro-batch within 15 minutes plus load time, and streaming within seconds to a minute.
The Core Methods of Real-Time Data Ingestion
Where the data originates decides the method.
Log-based CDC. This reads the database's own transaction log, the record it already keeps for crash recovery and replication:
- Postgres: the write-ahead log (WAL), read through a logical replication slot
- MySQL: the binary log, which must use row format
- MongoDB: change streams
CDC captures inserts, updates, and deletes, including hard deletes, without querying production tables. Coverage depends on configuration. In Postgres, only tables in the publication are captured, its publish setting controls which operations are emitted, and update and delete events identify rows by replica identity, the primary key by default. A published table with no usable identity rejects updates and deletes on the source. The source must also keep its log until the reader catches up, and Postgres needs wal_level = logical. Use it when a production database feeds a warehouse and you need low-latency, complete change capture.
Query-based polling. This runs SELECT * FROM orders WHERE updated_at > :last_synced_at on an interval. It needs no special setup, but it misses hard deletes, loses intermediate states between polls, and adds query load. The updated_at watermark is also fragile: a transaction that commits late with an earlier timestamp is skipped, and any writer that forgets to update the column goes unseen.
Trigger-based capture. Database triggers write each change to an audit table that another process reads. It catches deletes but adds write overhead to every transaction.
Event and API ingestion. Applications publish events, such as courier GPS pings, to a broker like Apache Kafka or Amazon Kinesis, or send them over an HTTP API. For event-native data such as telemetry, there may be no database row or transaction log to replicate. Applications can also publish business events derived from database transactions, often through an outbox pattern.
Most data ingestion tools specialize in one method. Debezium handles log-based CDC and publishes changes to Kafka. Managed platforms like Artie support CDC and API event ingestion.
Real-Time Data Ingestion: Architecture, Methods and How It Works
A data ingestion architecture for the delivery app's orders table, which lives in a Postgres database in us-east-1, has up to six components, one of them optional in simpler setups. The source is the database itself. The other five are below.
Capture. A CDC reader holds a logical replication slot and decodes the WAL into row-level change events: the operation, row identity when configured, and row data. The slot tells Postgres to keep WAL until the reader confirms it.
Transport (optional). Many topologies put a durable log, such as Kafka, Redpanda, or Kinesis, between capture and write. It absorbs temporary spikes, decouples capture from destination writes, and provides a durable replay buffer. It also preserves order within a partition, so keying by primary key keeps each row's changes in sequence. Simpler topologies skip it: source log retention plus durable checkpoints of the reader's position provide the recovery boundary. Our guide to real-time data streaming architecture covers these layers in more depth.
Processing. This optional layer filters columns, hashes PII, or joins reference data. Apache Flink handles stateful work like windowed aggregations. Many teams skip it and transform in the warehouse.
Destination writer. The writer applies change events to Snowflake as upserts and deletes keyed on the primary key. Delivery is at-least-once, so an event can arrive twice, and a MERGE on the primary key makes a replay produce the same final row. The pipeline must also preserve or reconcile source order for each key. Otherwise an older retry can overwrite newer state.
Observability. The pipeline reports lag from source commit to destination write, throughput, errors, and replication slot size.
How to Know If Your Workload Actually Needs Real-Time Ingestion
For each table you ingest, ask what breaks if the data is a day old. Real-time ingestion earns its cost when:
- A person or system acts on the data within minutes, like the ops dashboard
- Stale data causes direct loss, such as oversold inventory or an approved fraudulent charge
- Customer-facing features read from the warehouse or lake
Batch is enough when reports are read daily or weekly, finance reconciles at month-end, or training sets are built on a schedule. If the data is consumed less often than the batch runs, faster ingestion only adds cost: always-on compute, more moving parts, and someone on call. Many teams stream a few operational tables, like orders, and leave the rest on batch.
The Operational Challenges of Running Real-Time Ingestion at Scale
WAL growth on the source. A replication slot makes Postgres keep every WAL segment the reader hasn't confirmed. If the reader stalls, WAL piles up on the production primary until the disk fills. The max_slot_wal_keep_size setting is the safety cap, and it defaults to -1 (unlimited). This check shows each slot's state. wal_status needs Postgres 13 or later:
SELECT slot_name,
active,
wal_status,
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained_wal
FROM pg_replication_slots; slot_name | active | wal_status | retained_wal
--------------+--------+------------+--------------
orders_ingest | f | reserved | 38 GBA cap protects the primary by allowing required WAL to be removed when a slot falls too far behind. The slot can become unusable, and recovery commonly requires a new snapshot or backfill. Our replication slot guide covers the tradeoff.
Schema changes. A developer adds a courier_tip column to orders. Postgres logical replication doesn't carry DDL statements, so the pipeline has to detect the new column another way and add it to the destination before writing rows that use it.
Backfills. Existing rows have to land alongside live changes. Snapshotting a large table while updates keep arriving needs a coordinated consistency boundary: the snapshot reads as of a recorded log position and streaming resumes from that same position, so no change is missed or applied twice. A read replica can reduce snapshot load on the primary, but the snapshot boundary and CDC position still have to be coordinated.
Duplicates and ordering. Retries produce duplicates, so writes must be idempotent. Events for one order must also stay in order, or a late preparing event can overwrite delivered. In Kafka, ordering holds within a partition, so events are keyed by primary key.
TOAST columns. Postgres stores large values out of line and leaves unchanged ones out of update events. A writer that ignores this can silently corrupt the destination. We wrote up the failure in Why TOAST Columns Break Postgres CDC.
Teams solve all of these with Debezium, Kafka, and custom writers. Artie runs this layer as a managed service, with sub-minute latency for supported managed CDC workloads. Observed freshness depends on the source, destination, table mode, and workload. Backfills and traffic spikes can still add delay. Artie's founders started from warehouse data lagging production by hours to days, as described in Introducing Artie Transfer.
To act on this, pick the one table where staleness costs you the most and ingest only that. If you'd rather not run the capture, transport, and write layers yourself, start for free or talk to the Artie team.
FAQ
What is the difference between real-time ingestion and micro-batch ingestion?
Micro-batch ingestion collects records into small windows, often seconds to minutes long, and processes each window as a group. Real-time ingestion captures and delivers changes continuously, though the destination write may still be a short batch. The practical difference is whether a trigger interval sets your minimum latency.
Does real-time data ingestion require Kafka?
No. Kafka is a common transport layer because it buffers changes durably and preserves order within a partition. When a row's primary key is the partition key, that gives the pipeline a per-row ordering boundary. Simpler topologies skip it, relying on source-log retention plus durable checkpoints for recovery.
What is change data capture and how does it relate to real-time ingestion?
Change data capture (CDC) reads a database's transaction log and emits each insert, update, and delete as an event. It's a common method for real-time ingestion from production databases because it captures deletes and avoids querying production tables. Coverage depends on which tables and operations the source is configured to publish.
How does real-time ingestion handle schema changes in source databases?
The pipeline has to detect the change and apply it to the destination before writing rows that use it. Postgres logical replication doesn't carry DDL statements, so tools detect new columns another way. Handling varies by tool, so test column additions, drops, and type changes first.
What destinations can real-time ingestion pipelines write to?
Common destinations include cloud warehouses such as Snowflake, BigQuery, and Redshift, lakehouse platforms like Databricks, and message brokers like Kafka. Support varies by tool, so check the connector list for your exact source and destination pair before you commit to one.


