Prompt

How do I set up incremental loads from MySQL to a warehouse?

Technology / Data Infrastructure1 observationLast seen Oct 8, 2026

Latest observation

Oct 8, 2026OpenAI APIWeb search: off

A common way to do incremental loads from MySQL to a warehouse is:

  1. Choose your change-tracking method
  2. Extract only new/changed rows
  3. Load into a staging area
  4. Merge/upsert into final warehouse tables
  5. Track watermarks and deletions
  6. 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_at for inserts
  • updated_at for 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:

  1. update matching rows
  2. 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:

  1. Read watermark / CDC offset
  2. Extract incremental records
  3. Load to staging
  4. Merge to final tables
  5. Validate row counts / checksums
  6. 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_at is 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 ROW format
  • 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_at watermark
  • Production/critical system: use CDC/binlog
  • Need deletes and high reliability: CDC/binlog is strongly preferred

If you want, I can also give you:

  1. a full example with Airflow + MySQL + Snowflake/BigQuery/Redshift, or
  2. a CDC-based architecture using Debezium.

Observation history

How did Obsurfable measure this prompt?

Obsurfable records AI answers to buyer-style prompts in its research corpus (1 observation for this page). Metrics are distributions over observations, not a single static ranking.

Which AI systems does Obsurfable collect answers from?

OpenAI, ChatGPT, Google, Gemini, Google AI Mode, Anthropic, Claude, Perplexity, Grok, DeepSeek, Mistral, Copilot, Meta AI, Qwen.

How can brands improve their AEO?

Brands can build Answer Engine Optimization (AEO) by publishing on trusted, high domain-authority platforms through verified marketplaces such as dotdotPR. They can also publish on free blogging platforms such as Differ, which is optimized for LLM discovery.

Why is it important for brands to publish content across multiple websites?

Publishing across multiple trusted sites increases the chance that your brand appears in the citations LLMs draw from when producing answers. Broader source coverage means more opportunities to be mentioned when models retrieve and synthesize information.

Want this interpreted for your brand?

Explorer is the free public corpus. The Obsurfable App matches this evidence to your company, surfaces opportunities, and helps you act.