为什么全量刷新不够用

在 ETL 入门阶段,全量刷新听起来很省事:每次把源表的所有数据删掉,重新拉一遍。对于小数据量表,这没什么问题。但当数据量增长到百万、千万甚至亿级时,全量刷新就成了一场噩梦。

来看几个数字。假设你有一个订单表,每天新增 5 万条记录,总计 2000 万行。全量读取 2000 万行可能耗费 30 分钟,但实际发生变化的数据只有 5 万行,占总量的 0.25%。这意味着你浪费了 99.75% 的传输量和计算量。

全量刷新还有几个更棘手的问题:

  • 窗口时间不够:数据量越大,刷新窗口越长。如果每天只有 2 小时的维护窗口,全量同步 10 亿条数据根本跑不完
  • 对源库压力大:全表扫描会对业务数据库造成巨大压力,影响在线交易
  • 浪费存储和带宽:每次传输全部数据,压缩、网络、磁盘 IO 的成本成倍增加

增量抽取就是为了解决这些问题而生的。

增量抽取的三种主要策略

增量抽取的核心是回答一个问题:“上次处理完之后,哪些数据发生了变化?” 针对不同的数据源和业务场景,有三种主流策略。

时间戳增量

时间戳策略是最简单、应用最广泛的增量抽取方式。原理是:源表里有一个时间字段(如 updated_atmodified_date),记录每条数据最后修改的时间。每次抽取时,只取时间大于上次抽取时间的记录。

-- 最后一次抽取的时间是 2026-04-12 03:00:00
-- 只抽取此后修改的数据
SELECT * FROM orders
WHERE updated_at > '2026-04-12 03:00:00'
  AND updated_at <= '2026-04-13 03:00:00';

时间戳策略的优点是实现简单,对源库改造小。但缺点也很明显:

  • 需要源表有时间戳字段,且业务代码正确维护该字段
  • 物理删除的数据无法被捕获(删除不更新时间戳)
  • 时间精度不够时可能漏数据(如果两行在同一毫秒内修改)
  • 批量更新时间戳的情况可能导致重复或遗漏

增量文件/日志

有些系统会主动生成包含变更的文件或日志。例如:

  • 应用日志:业务系统把每次增删改操作输出到日志文件
  • 增量文件:ERP 系统每天生成变更文件,记录当天修改的记录 ID
  • Binlog:MySQL 的二进制日志记录了所有数据变更

增量文件的方式对源系统侵入最小,但依赖应用配合生成这些文件。

基于日志的 CDC(Change Data Capture)

CDC 是增量抽取的终极方案。它的原理是直接读取数据库的事务日志(如 MySQL 的 binlog、PostgreSQL 的 WAL),从中解析出所有数据变更事件。

-- 查看 MySQL 是否开启 binlog
SHOW VARIABLES LIKE 'log_bin';
SHOW VARIABLES LIKE 'binlog_format';

-- 确保 binlog 格式为 ROW
-- ROW 格式记录每行数据变更前后的值,是最适合 CDC 的模式

基于日志的 CDC 有几个压倒性的优势:

  • 零侵入:不需要修改业务代码,不需要在表上加时间戳
  • 捕获所有变更:包括 INSERT、UPDATE、DELETE,连 DDL 都能捕获
  • 近乎实时:数据一旦提交,CDC 工具马上就能感知到
  • 一致性强:直接从事务日志读,保证不丢数据

CDC 三种方法详解

从实现层面看,CDC 有三种具体的方法,我们逐个分析。

数据库触发器 CDC

在源数据库上建立触发器,每当有 INSERT、UPDATE、DELETE 操作时,触发器把变更记录写入一张专门的审计表(Changelog Table)。

-- 创建变更记录表
CREATE TABLE orders_changelog (
    log_id BIGINT AUTO_INCREMENT PRIMARY KEY,
    table_name VARCHAR(100),
    operation_type ENUM('INSERT', 'UPDATE', 'DELETE'),
    record_id BIGINT,
    old_data JSON,
    new_data JSON,
    changed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    processed BOOLEAN DEFAULT FALSE
);

-- 创建触发器捕获 INSERT
CREATE TRIGGER trg_orders_insert
AFTER INSERT ON orders
FOR EACH ROW
INSERT INTO orders_changelog (table_name, operation_type, record_id, new_data)
VALUES ('orders', 'INSERT', NEW.id, JSON_OBJECT(
    'order_id', NEW.order_id,
    'amount', NEW.amount,
    'status', NEW.status
));

-- 创建触发器捕获 UPDATE
CREATE TRIGGER trg_orders_update
AFTER UPDATE ON orders
FOR EACH ROW
INSERT INTO orders_changelog (table_name, operation_type, record_id, old_data, new_data)
VALUES ('orders', 'UPDATE', NEW.id,
    JSON_OBJECT(
        'order_id', OLD.order_id,
        'amount', OLD.amount,
        'status', OLD.status
    ),
    JSON_OBJECT(
        'order_id', NEW.order_id,
        'amount', NEW.amount,
        'status', NEW.status
    )
);

触发器方案的优点是全兼容,任何支持触发器的数据库都能用。缺点是对源库性能有影响——每次 DML 操作都额外多写一次变更表。在高并发场景下,触发器可能成为瓶颈。

事务日志挖掘 CDC(Debezium)

这是目前最主流的 CDC 方式。Debezium 是一个基于 Kafka Connect 的开源 CDC 平台,支持 MySQL、PostgreSQL、MongoDB、SQL Server 等多种数据库。

Debezium 的工作原理如下:

  1. 连接到数据库,读取 binlog 或 WAL 日志
  2. 解析日志中的变更事件
  3. 把变更事件写入 Kafka 主题
  4. 下游消费者从 Kafka 读取并处理变更

一个 Debezium 捕获的变更事件结构如下:

{
  "schema": { "...": "..." },
  "payload": {
    "op": "u",
    "ts_ms": 1712880000123,
    "before": {
      "order_id": 1001,
      "amount": 299.00,
      "status": "PAID"
    },
    "after": {
      "order_id": 1001,
      "amount": 299.00,
      "status": "SHIPPED"
    },
    "source": {
      "db": "ecommerce",
      "table": "orders",
      "ts_ms": 1712880000100
    }
  }
}

在 Debezium 中,每条变更记录包含操作类型(op)、变更前的数据(before)、变更后的数据(after)以及源信息(source)。其中 op 字段的含义是:

  • c:创建(Create),对应 INSERT
  • u:更新(Update),对应 UPDATE
  • d:删除(Delete),对应 DELETE
  • r:读取(Read),用于快照

使用 Docker 启动 Debezium 连接 MySQL 的步骤非常简单:

version: '3.8'
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.6.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  kafka:
    image: confluentinc/cp-kafka:7.6.0
    ports:
      - "9092:9092"
    depends_on:
      - zookeeper
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1

  debezium-connector:
    image: debezium/connect:2.5
    ports:
      - "8083:8083"
    depends_on:
      - kafka
    environment:
      BOOTSTRAP_SERVERS: kafka:9092
      GROUP_ID: 1
      CONFIG_STORAGE_TOPIC: connect-configs
      OFFSET_STORAGE_TOPIC: connect-offsets
      STATUS_STORAGE_TOPIC: connect-status

启动服务后,注册一个 Debezium 连接器:

curl -i -X POST localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d '{
    "name": "orders-connector",
    "config": {
      "connector.class": "io.debezium.connector.mysql.MySqlConnector",
      "database.hostname": "mysql-host",
      "database.port": "3306",
      "database.user": "debezium",
      "database.password": "dbz_pass",
      "database.server.id": "184054",
      "database.server.name": "ecommerce",
      "database.include.list": "ecommerce",
      "table.include.list": "ecommerce.orders",
      "database.history.kafka.bootstrap.servers": "kafka:9092",
      "database.history.kafka.topic": "schema-changes.orders",
      "include.schema.changes": "true",
      "snapshot.mode": "initial"
    }
  }'

注册完成后,Debezium 会自动读取 MySQL 的 binlog,把 orders 表的每一次变更推送到 Kafka 主题 ecommerce.ecommerce.orders。其他系统只需要从这个主题消费就能获取实时的数据变更。

Debezium 是目前最流行的 CDC 工具之一。它的优势在于与 Kafka 深度集成,天然支持消息持久化和回溯消费。但缺点是运维复杂度较高,需要同时维护 Kafka 和 Kafka Connect 集群。

查询式 CDC

查询式 CDC 不做日志解析,而是通过 SQL 查询来找到变更数据。常见做法是在源表上维护水印字段,每次查询大于水印值的数据。

import psycopg2
from datetime import datetime


class QueryBasedCDC:
    """基于查询的 CDC 实现"""
    
    def __init__(self, conn_params, watermark_table="etl_watermarks"):
        self.conn = psycopg2.connect(**conn_params)
        self.watermark_table = watermark_table
        self._init_watermark_table()
    
    def _init_watermark_table(self):
        """初始化水印表"""
        with self.conn.cursor() as cur:
            cur.execute(f"""
                CREATE TABLE IF NOT EXISTS {self.watermark_table} (
                    source_table VARCHAR(255) PRIMARY KEY,
                    last_max_id BIGINT DEFAULT 0,
                    last_max_timestamp TIMESTAMP,
                    updated_at TIMESTAMP DEFAULT NOW()
                )
            """)
            self.conn.commit()
    
    def get_watermark(self, table_name):
        """获取水印值"""
        with self.conn.cursor() as cur:
            cur.execute(
                f"SELECT last_max_id FROM {self.watermark_table} WHERE source_table = %s",
                (table_name,)
            )
            row = cur.fetchone()
            return row[0] if row else 0
    
    def update_watermark(self, table_name, max_id):
        """更新水印值"""
        with self.conn.cursor() as cur:
            cur.execute(f"""
                INSERT INTO {self.watermark_table} (source_table, last_max_id, updated_at)
                VALUES (%s, %s, NOW())
                ON CONFLICT (source_table)
                DO UPDATE SET last_max_id = %s, updated_at = NOW()
            """, (table_name, max_id, max_id))
            self.conn.commit()
    
    def fetch_incremental(self, source_table, batch_size=5000):
        """增量获取数据"""
        last_id = self.get_watermark(source_table)
        
        with self.conn.cursor() as cur:
            cur.execute(f"""
                SELECT * FROM {source_table}
                WHERE id > %s
                ORDER BY id
                LIMIT %s
            """, (last_id, batch_size))
            
            columns = [desc[0] for desc in cur.description]
            rows = cur.fetchall()
            
            if not rows:
                return []
            
            max_id = rows[-1][0]
            
            result = [dict(zip(columns, row)) for row in rows]
            self.update_watermark(source_table, max_id)
            
            return result


# 使用示例
cdc = QueryBasedCDC({
    "host": "localhost",
    "port": 5432,
    "dbname": "ecommerce",
    "user": "etl_user",
    "password": "etl_pass"
})

while True:
    batch = cdc.fetch_incremental("orders", batch_size=5000)
    if not batch:
        print("没有新数据,等待 30 秒...")
        time.sleep(30)
        continue
    
    print(f"获取到 {len(batch)} 条新订单")
    process_orders(batch)

查询式 CDC 的优点是实现简单,不需要额外的基础设施。缺点是延迟较高(依赖定时轮询),并且无法捕获物理删除。

CDC 工具的对比

在实际项目中,选择合适的 CDC 工具需要考虑多个维度。下表对比了三种常见的工具选型:

特性 Debezium AWS DMS Fivetran
部署方式 自托管 托管服务 SaaS
支持的源库 MySQL, PG, SQL Server, MongoDB 等 20+ 种 150+ 种
目标端 Kafka 多种数据库/S3 数据仓库
延迟 秒级 秒级 分钟级
运维复杂度
成本 开源免费 按实例计费 按行数计费
数据转换 需自建 可配置映射 内置 schema 迁移

选型建议:

  • 如果团队有 Kafka 运维能力,Debezium 是最灵活的选择
  • 如果不想操心基础设施,AWS DMS 或 Fivetran 这类托管服务更省心
  • 如果是简单的增量抽取场景,查询式 CDC + 水印表就足够了

处理延迟到达的数据

现实世界中的数据不是完美的。由于网络延迟、批处理排队、上游系统故障等原因,部分数据可能会延迟到达。具体来说:

  • 某笔订单昨晚 10 点创建,但今早 8 点才到达 ETL 系统
  • 上游系统补录了一条昨天的数据,更新时间戳是今天

处理延迟数据有几种策略:

基于事件时间的处理

不依赖数据到达时间,而是依赖数据本身的业务时间(如 order_created_at)。这种方式在 Flink 和 Spark Streaming 中有成熟的支持。

def process_with_event_time(record):
    """基于事件时间处理"""
    # 使用业务时间而不是系统时间
    business_date = record["order_created_at"].date()
    target_partition = f"orders_{business_date.strftime('%Y%m%d')}"
    write_to_partition(record, target_partition)

允许更新的水印策略

水印不应该是单向增大的。允许一定的时间窗口内水印回退或重新处理某个时间段的数据。

class SlidingWatermark:
    """滑动窗口水印管理"""
    
    def __init__(self, lookback_hours=24):
        self.watermark_table = "etl_watermarks"
        self.lookback = timedelta(hours=lookback_hours)
    
    def get_safe_window(self):
        """获取安全处理窗口"""
        last_watermark = self.read_watermark()
        window_start = last_watermark - self.lookback
        return window_start, last_watermark
    
    def process_late_data(self, late_record):
        """处理延迟到达的数据"""
        event_time = late_record["event_time"]
        safe_start, safe_end = self.get_safe_window()
        
        if safe_start <= event_time <= safe_end:
            # 在安全窗口内,直接写入
            self.upsert_to_target(late_record)
        elif event_time < safe_start:
            # 太旧的数据,需要人工确认
            self.send_to_dead_letter(late_record, "数据超出回溯窗口")

精确一次处理语义

在 CDC 场景中,最糟糕的结果不是数据重复,而是数据丢失。重复数据可以通过去重来纠正,丢失的数据则很难发现。

实现精确一次处理需要三个层面的保证:

源端:至少一次捕获

CDC 工具需要记录已经读取到的日志位置(offset)。常见做法是把 offset 存到 Kafka 或独立的元数据存储中。

def capture_with_offset(debezium_record, offset_storage):
    """带偏移量记录的 CDC 消费"""
    offset = debezium_record["source"]["offset"]
    
    # 检查是否处理过
    if offset_storage.is_processed(offset):
        print("这条变更已经处理过,跳过")
        return
    
    # 记录偏移量
    offset_storage.mark_processing(offset)
    
    try:
        # 处理变更
        process_change(debezium_record)
        offset_storage.mark_completed(offset)
    except Exception as e:
        print(f"处理失败: {e}")
        offset_storage.mark_failed(offset)

目标端:幂等写入

前面文章提过,幂等性是批处理的核心原则。在 CDC 场景中同样适用。最好的做法是使用 UPSERT(MERGE)来写入目标表。

def upsert_to_warehouse(records, target_table):
    """幂等写入目标表"""
    # 使用 MERGE 语句避免重复
    merge_sql = f"""
    MERGE INTO {target_table} AS target
    USING (SELECT %s AS id, %s AS amount, %s AS status) AS source
    ON target.id = source.id
    WHEN MATCHED THEN
        UPDATE SET amount = source.amount, status = source.status
    WHEN NOT MATCHED THEN
        INSERT (id, amount, status) VALUES (source.id, source.amount, source.status)
    """
    
    for record in records:
        execute(merge_sql, (record["id"], record["amount"], record["status"]))

协调:两阶段提交或事务

对于要求极高一致性的场景,可以结合事务性消息和两阶段提交。Kafka 支持事务性写入,配合下游的事务性读取可以实现端到端精确一次。

# Kafka 事务性生产
producer.init_transactions()
producer.begin_transaction()

try:
    for record in changes:
        producer.send("cdc-orders", record)
    producer.commit_transaction()
except Exception:
    producer.abort_transaction()

小结

增量抽取和 CDC 是解决大规模数据同步问题的核心手段。从简单的时间戳策略到基于日志的 Debezium,每种方法都有自己的适用场景。选择 CDC 方案时要考虑数据量大小、实时性要求、源库类型和团队运维能力。延迟到达的数据要采用事件时间处理或滑动窗口来应对,而精确一次处理依赖源头 offset 记录、目标端幂等写入、以及全链路的容错设计。

Summary: 增量抽取的三种策略和 CDC 原理,涵盖时间戳、日志、查询式方案及精确一次处理。