Data ingestion is the first technical layer of any pipeline. Pick the wrong pattern and you build operational pain forever.
Pattern 1 — Full table snapshot (batch)
Every run: read the entire source table; load into the warehouse.
Run nightly:
SELECT * FROM source.orders
→ load to warehouse.raw.orders (full replace)
Pros:
- Simple. No state to track.
- Works for any source that can do SELECT *.
- Self-healing (each load is a full picture).
Cons:
- Expensive for large tables (re-reads everything every run).
- No history of changes (every load overwrites).
- Latency = batch interval.
Best for:
- Small tables (< 1M rows).
- Reference data that doesn't change often.
- When you don't care about historical state.
Pattern 2 — Incremental batch (high-water mark)
Each run: read only rows changed since last run, identified by a watermark column.
Run hourly:
SELECT * FROM source.events WHERE created_at > <last_run_timestamp>
→ append to warehouse.raw.events
Pros:
- Much cheaper than full snapshot.
- Scales to large tables.
- Preserves history (rows accumulate).
Cons:
- Requires a reliable watermark column (
updated_atorcreated_at). - Late-arriving data (events synced after the watermark) get dropped.
- Watermark management is stateful.
Best for:
- Append-mostly tables (events, transactions, logs).
- Large tables where full snapshot is expensive.
Use sync time (_loaded_at) rather than event time when possible to avoid late-arrival drops.
Pattern 3 — CDC (Change Data Capture)
The source's transaction log is streamed; every INSERT, UPDATE, DELETE is captured and replicated.
Postgres WAL or MySQL binlog → CDC tool → warehouse
Pros:
- Near-real-time (seconds to minutes lag).
- Captures DELETEs and UPDATEs correctly.
- No load on source DB (reads the log, not the tables).
- Catches every change, including very late updates.
Cons:
- Requires DB-level access (log reader).
- More operational complexity.
- CDC tools cost money (Fivetran, Debezium, etc.).
- Tricky to backfill historical state (you typically snapshot + start CDC from a point).
Best for:
- OLTP databases where you need close-to-real-time and accurate.
- Tables with frequent UPDATEs (incremental batch misses them).
- Compliance / audit use cases.
Major CDC tools:
- Fivetran — managed, expensive.
- Airbyte CDC — open source, self-hosted.
- Debezium — open source, requires Kafka infrastructure.
- AWS DMS — AWS-native.
Pattern 4 — Streaming (event-time)
Application emits events to a stream (Kafka, Kinesis, Pub/Sub); consumer reads stream and loads warehouse.
App → Kafka topic → stream consumer → warehouse table
Pros:
- True real-time when needed.
- Decouples producer from consumer.
- Replayable (events stay in stream).
Cons:
- Infrastructure complexity (Kafka cluster or managed equivalent).
- Stream processor needed (Flink, Spark Streaming, Materialize, RisingWave).
- Out-of-order events, exactly-once semantics — operational concerns.
Best for:
- Genuinely real-time use cases (fraud, alerting, live dashboards).
- High-volume event data where batch struggles.
- Decoupled architectures with many consumers.
Pattern 5 — File-based
Source produces files (CSV, JSON, Parquet) on S3/GCS/SFTP; pipeline picks them up.
External vendor uploads daily file to s3://incoming/
→ Lambda / Airflow detects new file → loads to warehouse
Pros:
- Simple. No source DB access needed.
- Works with vendors who don't offer API.
- Easy backfill (re-process files).
Cons:
- Latency depends on file drop schedule.
- Schema can drift unexpectedly.
Best for:
- B2B vendor data feeds.
- Marketing platform exports.
- Compliance/audit logs.
Decision tree
Where does the data live?
├─ External SaaS (Salesforce, Stripe, etc.) → Fivetran/Airbyte (uses their API; CDC where supported)
├─ Your OLTP database (Postgres/MySQL)
│ ├─ Small + slow-changing → full snapshot batch
│ ├─ Large + append-mostly → incremental batch
│ └─ Need correctness on UPDATEs/DELETEs OR near-real-time → CDC
├─ Application events → Kafka/streaming (if real-time) OR batch event load
└─ Files from a vendor → file-watch pipeline
Cost reality
For 100K rows/hour:
| Pattern | Cost | Latency |
|---|---|---|
| Full snapshot (hourly) | Modest, scales with table size | 1h |
| Incremental batch | Low | 1h |
| CDC | Modest (tool cost) | 1-15 min |
| Streaming (own infra) | High (Kafka + ops) | Seconds |
| Managed connector (Fivetran) | High (per-row fees) | 5-30 min |
Build vs buy
For each source, choose:
- Managed (Fivetran/Airbyte Cloud) — pay per row; saves engineer hours.
- Self-hosted (Airbyte OSS, Meltano) — operate yourself; saves money at scale.
- Custom Python script — build for unusual sources.
Most teams use Fivetran for 5-10 critical sources and skip it for one-off custom ones.
Threshold: Fivetran costs 5-10x what it costs to operate Airbyte OSS at scale. For high-volume connectors, switching to OSS pays back in 6-12 months.
Schema evolution
Sources change. Columns added, removed, renamed.
How each pattern handles it:
- Full snapshot: detects schema each run. Easy to handle additions; deletions/renames more disruptive.
- Incremental: same as snapshot.
- CDC: tools often pause on schema changes; manual intervention.
- Streaming: schema registry (Confluent Schema Registry, Apicurio) helps.
Plan for schema evolution. It will happen.
Common ingestion mistakes
- Full snapshot of huge tables. Bandwidth and warehouse load mount.
- Incremental on event_time without lookback. Late-arriving events lost.
- CDC without monitoring. Lag grows silently when consumers crash.
- Streaming for non-real-time use cases. Complexity tax for no benefit.
- No dedup logic for "at-least-once" deliveries. Same row appears twice. Add
_synced_at + ROW_NUMBER()dedup in staging.
Takeaway
Match pattern to source. Snapshot for small/static. Incremental for append-mostly. CDC for OLTP needing accuracy or low-latency. Streaming for genuine real-time. Build/buy is per-source; managed tools save time, OSS saves money at scale.
📘 Companion Deep Dive: Production Python Ingestion Patterns
While managed ELT tools (Airbyte, Fivetran) and CDC cover database replication, data engineers routinely write custom Python extractors for internal microservices and third-party SaaS APIs.
1. The Pitfall of Offset-Based Pagination
Querying GET /v1/transactions?limit=100&offset=5000 causes database query performance to degrade quadratically ($O(N^2)$) on the server. Furthermore, if new records are inserted while pagination is underway, records shift forward, producing duplicated or skipped events.
2. Production Cursor-Based Pagination
Always favor cursor-based (keyset) pagination using monotonic IDs or timestamps:
import requests
import time
from typing import Generator, Dict, Any
def fetch_paginated_api(base_url: str, api_key: str) -> Generator[Dict[str, Any], None, None]:
headers = {"Authorization": f"Bearer {api_key}", "Accept": "application/json"}
cursor = None
has_more = True
while has_more:
params = {"limit": 100}
if cursor:
params["starting_after"] = cursor
# Exponential backoff on HTTP 429 / 5xx
for attempt in range(5):
response = requests.get(base_url, headers=headers, params=params, timeout=30)
if response.status_code == 200:
break
elif response.status_code == 429 or response.status_code >= 500:
sleep_time = (2 ** attempt) + (time.time() % 1)
time.sleep(sleep_time)
else:
response.raise_for_status()
payload = response.json()
data = payload.get("data", [])
if not data:
break
for record in data:
yield record
cursor = data[-1].get("id")
has_more = payload.get("has_more", False)
3. Raw Landing Envelopes
Never land bare payload records directly into storage without audit metadata. Wrap every payload in a standardized metadata envelope:
{
"_raw_payload": { "id": "tx_99812", "amount": 4200, "status": "settled" },
"_ingested_at": "2026-09-16T14:30:00Z",
"_source_system": "billing_service_prod",
"_pipeline_run_id": "run_01j7x4v9m"
}