Prompt
How do I set up incremental loads from MySQL to a warehouse?
Latest observation
A common way to do incremental loads from MySQL to a warehouse is:
- Choose your change-tracking method
- Extract only new/changed rows
- Load into a staging area
- Merge/upsert into final warehouse tables
- Track watermarks and deletions
- Automate and monitor
Here’s a practical setup.
1) Pick the incremental strategy
Option A: updated_at / created_at watermark
Best if your tables have reliable timestamp columns.
created_atfor insertsupdated_atfor updates- Your job queries rows where
updated_at > last_successful_run
Pros: simple
Cons: misses deletes unless you handle them separately; timestamps must be trustworthy
Option B: MySQL binlog / CDC
Use change data capture tools like:
- Debezium
- Fivetran
- Airbyte
- AWS DMS
- Maxwell
These read MySQL’s binary logs and capture inserts, updates, deletes.
Pros: captures deletes and updates reliably
Cons: more setup, usually needs infrastructure and permissions
If you want a robust warehouse pipeline, CDC/binlog is usually the best long-term approach.
2) Set up a source tracking mechanism
If using timestamps
Store a watermark table in your orchestration/db metadata.
Example:
CREATE TABLE etl_watermarks (
pipeline_name VARCHAR(100) PRIMARY KEY,
last_run_at DATETIME NOT NULL
);
At the end of a successful run, update last_run_at.
Query example:
SELECT *
FROM orders
WHERE updated_at > '2026-10-08 10:00:00'
AND updated_at <= '2026-10-08 11:00:00';
Use a small overlap window if needed to avoid missing late-arriving updates.
If using CDC/binlog
Track offsets automatically via your CDC tool:
- binlog file + position
- GTID
- event sequence
No manual watermarking needed in most cases.
3) Extract into a staging table
Load incremental data into a raw/staging table first, not directly into final analytics tables.
Example staging table:
CREATE TABLE stg_orders (
order_id BIGINT,
customer_id BIGINT,
amount DECIMAL(10,2),
updated_at DATETIME,
_load_ts DATETIME
);
This gives you:
- auditability
- reprocessing capability
- validation before merge
4) Merge into the warehouse table
Your warehouse usually needs an upsert/merge step.
Example pattern:
MERGE INTO orders AS target
USING stg_orders AS source
ON target.order_id = source.order_id
WHEN MATCHED AND source.updated_at > target.updated_at THEN
UPDATE SET
customer_id = source.customer_id,
amount = source.amount,
updated_at = source.updated_at
WHEN NOT MATCHED THEN
INSERT (order_id, customer_id, amount, updated_at)
VALUES (source.order_id, source.customer_id, source.amount, source.updated_at);
If your warehouse doesn’t support MERGE, do:
- update matching rows
- insert new rows
5) Handle deletes
Deletes are the tricky part.
If using timestamp watermark only
You won’t see deletes unless:
- you have a soft-delete column like
is_deleted - you maintain delete logs separately
- you periodically reconcile full snapshots
If using CDC/binlog
Deletes come through naturally.
In the warehouse, you can:
- hard delete
- or soft delete (
is_deleted = true)
Soft delete is often safer for analytics.
6) Make loads idempotent
Idempotency means rerunning a job won’t duplicate or corrupt data.
Best practices:
- Load to staging first
- Deduplicate by primary key + latest
updated_at - Use
MERGE/upsert logic - Update watermark only after successful load
Example dedupe step:
SELECT *
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY updated_at DESC) AS rn
FROM stg_orders
) x
WHERE rn = 1;
7) Orchestrate the pipeline
Use a scheduler/orchestrator like:
- Airflow
- Dagster
- Prefect
- dbt Cloud
- cron for simple cases
Typical flow:
- Read watermark / CDC offset
- Extract incremental records
- Load to staging
- Merge to final tables
- Validate row counts / checksums
- Advance watermark / offset
8) Add data quality checks
At minimum check:
- row counts extracted vs loaded
- no duplicate primary keys in staging
- nulls in required columns
- freshness SLA
- late-arriving data behavior
Example:
- alert if no rows loaded for a table that normally changes daily
- compare source and target counts for the batch
9) Recommended architecture
Simple setup
- MySQL table with
updated_at - Airflow scheduled job
- Extract rows since last watermark
- Load to warehouse staging
- Merge into target
Better production setup
- MySQL binlog CDC via Debezium/Fivetran/Airbyte/DMS
- Raw landing table
- Transform/merge in warehouse
- Monitor lag and failures
10) Important MySQL considerations
If you’re using timestamp-based increments:
- Ensure
updated_atis indexed - Use UTC to avoid timezone issues
- Handle clock skew carefully
- Don’t rely on
NOW()from different systems without standardization
If using CDC:
- Enable binlog in MySQL
- Prefer
ROWformat - Ensure binlog retention is long enough
- Handle schema changes explicitly
Example incremental SQL pattern
SELECT *
FROM orders
WHERE updated_at > :last_run_at
AND updated_at <= :current_run_at
ORDER BY updated_at;
Then in the warehouse:
MERGE INTO dw.orders t
USING stg.orders s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...;
Rule of thumb
- Small/simple system: use
updated_atwatermark - Production/critical system: use CDC/binlog
- Need deletes and high reliability: CDC/binlog is strongly preferred
If you want, I can also give you:
- a full example with Airflow + MySQL + Snowflake/BigQuery/Redshift, or
- a CDC-based architecture using Debezium.