数据库基础体系 · 第 137/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

ETL 与 ELT 数据管道:批流处理、幂等、补数、校验和可观测

数据管道的目标不是“把数据从 A 搬到 B”,而是在源系统、消息系统、计算引擎和目标存储之间,建立一条可重复执行、可恢复、可验证、可解释的数据处理链路。

一条成熟的数据管道至少要回答以下问题:

  • 数据从哪里来,读取边界是什么?
  • 转换发生在源端、管道中,还是目标端?
  • 一批数据或一条事件处理到哪一步了?
  • 任务失败后重跑,会不会重复写入或覆盖正确结果?
  • 延迟到达、乱序、重复消费和源端更新如何处理?
  • 历史数据修复或业务规则变更后,如何补数?
  • 如何证明目标数据没有少、没有多、没有被错误转换?
  • 发现问题后,能否定位到具体批次、分区、消息或 SQL?

这些问题分别对应 ETL 与 ELT、批流处理、幂等、补数、校验和可观测性。


一、先建立数据管道的基本模型

可以把数据管道抽象为:

SExtractRTransformDLoadTS \xrightarrow{\text{Extract}} R \xrightarrow{\text{Transform}} D \xrightarrow{\text{Load}} T

其中:

  • SS:源系统,例如 OLTP 数据库、文件、API、消息队列;
  • RR:原始数据或中间数据;
  • DD:经过清洗、关联、聚合后的数据;
  • TT:目标系统,例如数据仓库、湖仓、搜索索引或报表库。

实际系统通常还包含状态:

Pipeline=Data+State+Side Effects\text{Pipeline} = \text{Data} + \text{State} + \text{Side Effects}

状态包括:

  • 已读取到的源端位置;
  • 已消费的消息 offset;
  • 已完成的批次;
  • 已处理的业务主键;
  • 当前使用的规则或代码版本;
  • 失败记录和重试次数。

副作用包括:

  • 写入目标数据库;
  • 发送下游消息;
  • 更新搜索索引;
  • 触发告警或回调。

只保存“处理进度”而不保存“处理结果的身份”,通常不足以实现可靠恢复。因为进度提交和结果写入可能不是同一个事务。

例如,消费者执行:

  1. 读取消息 offset=100
  2. 写入目标数据库;
  3. 提交 offset。

如果第 2 步成功、第 3 步失败,重启后会再次读取 offset=100。因此系统必须允许第 2 步重复执行而结果不变,或者让 offset 与目标写入处于同一个可协调的事务边界内。


二、ETL 与 ELT:差别不只是字母顺序

2.1 ETL 的处理顺序

ETL 是:

  1. Extract:从源系统提取;
  2. Transform:在管道计算层或专用转换层处理;
  3. Load:将结果写入目标系统。

示意:

OLTP / 文件 / API
        |
        v
   提取与转换引擎
        |
        v
  数据仓库 / 报表库

ETL 常见于:

  • 目标系统计算能力有限;
  • 需要在进入目标前脱敏、过滤或格式转换;
  • 目标系统不适合保存原始数据;
  • 数据必须在进入共享区域前满足访问控制要求。

例如,源表包含身份证号和手机号,管道先将敏感字段脱敏,再写入分析库。

2.2 ELT 的处理顺序

ELT 是:

  1. Extract:提取原始数据;
  2. Load:先加载到目标存储;
  3. Transform:利用目标仓库或湖仓的计算能力转换。

示意:

OLTP / 文件 / 消息
        |
        v
  原始层 / Bronze
        |
        v
  清洗层 / Silver
        |
        v
  汇总层 / Gold

ELT 的重要特点是保留原始数据。这样可以在业务规则变化后重新转换,而不必再次访问源系统。

例如,订单原始事件先保存为:

{
  "event_id": "e-1001",
  "order_id": "o-10",
  "event_type": "paid",
  "occurred_at": "2025-01-01T10:00:00Z",
  "amount": 199.00
}

之后可以按不同业务规则生成:

  • 支付成功订单表;
  • 按日收入表;
  • 用户生命周期指标;
  • 财务对账表。

2.3 ETL 与 ELT 的边界

ETL 和 ELT 不是互斥架构。实际系统常见的是混合方式:

源库
  -> CDC 原始层
  -> 轻量清洗
  -> 湖仓明细层
  -> 仓库内 SQL 转换
  -> 数据集市

这里同时存在:

  • CDC 消息的轻量 ETL;
  • 原始数据进入湖仓的 Load;
  • 目标仓库内的 ELT。

选择依据主要是:

  • 源系统负载;
  • 目标系统的计算能力;
  • 原始数据是否需要保留;
  • 数据隐私与合规要求;
  • 数据延迟要求;
  • 转换逻辑是否需要重放;
  • 目标系统是否支持事务、MERGE、分区替换等操作。

不能简单认为“ELT 一定先进”或“ETL 一定更安全”。ELT 保存原始数据会增加存储和治理责任;ETL 提前丢弃字段则可能使后续修复无法进行。


三、批处理与流处理

3.1 批处理的边界

批处理把有限范围的数据作为一个集合处理。批次通常由以下任一种边界定义:

  • 时间区间:2025-01-01 00:00:00 <= t < 2025-01-02 00:00:00
  • 分区:某个日期分区或文件分区;
  • 自增 ID:100000 <= id < 110000
  • CDC 日志位置:某个 offset、LSN 或 binlog 位点;
  • 业务批次号:一次导入任务产生的一组文件。

批处理的关键不是“定时执行”,而是必须明确:

  1. 批次范围;
  2. 批次状态;
  3. 成功条件;
  4. 失败后如何重跑;
  5. 结果是否可以替换或合并。

例如,按事件时间处理 2025 年 1 月 1 日的数据,应使用半开区间:

[2025-01-01T00:00:00, 2025-01-02T00:00:00)[2025\text{-}01\text{-}01T00:00:00,\ 2025\text{-}01\text{-}02T00:00:00)

半开区间可以避免相邻批次重复或遗漏边界值。

3.2 流处理的边界

流处理面对的是持续到达的数据。它通常以以下状态推进:

  • 消息 offset;
  • Kafka 分区位点;
  • 数据库 CDC 的 LSN 或 binlog position;
  • 事件时间 watermark;
  • 窗口的完成状态。

流处理并不意味着“每条消息都立即得到最终结果”。只要允许乱序和迟到,系统就必须决定何时认为一个窗口暂时完成。

例如,统计每 5 分钟的支付金额:

Wk=[t0+5k, t0+5(k+1))W_k = [t_0 + 5k,\ t_0 + 5(k+1))

如果事件时间为 10:04:59 的消息在 10:08 才到达,则它属于 10:00–10:05 窗口,但窗口可能已经输出过一次。系统需要:

  • 丢弃迟到数据;
  • 更新窗口结果;
  • 重新计算整个窗口;
  • 将迟到数据放入修正队列。

这不是实现细节,而是业务语义。财务统计通常不能简单丢弃迟到数据;实时看板可能接受先出近似值、后续修正。

3.3 处理时间、事件时间与摄入时间

常见的三个时间含义不同:

  • 事件时间:业务事件实际发生的时间;
  • 处理时间:计算引擎处理该事件的时间;
  • 摄入时间:系统将事件写入管道或目标存储的时间。

例如:

订单在 10:00 创建
源库在 10:01 写入
CDC 在 10:02 捕获
消费者在 10:03 处理

这四个时刻不能混用。按处理时间统计可能导致跨日错分;按事件时间统计则要处理迟到、时钟漂移和修正。


四、CDC、消息顺序与重复消费

当 OLTP 数据库是源系统时,常见做法是读取数据库变更日志,而不是反复扫描业务表。CDC 可以捕获:

  • 插入;
  • 更新;
  • 删除;
  • 事务边界;
  • 日志位置;
  • 部分实现中的事务提交时间和表结构信息。

PostgreSQL 通常以 WAL 及其逻辑复制能力为基础;MySQL 通常以 binlog 为基础。Debezium 等 CDC 工具会把这些变更转换为消息。

但是,CDC 消息不是天然等价于最终业务事实。

4.1 更新事件与快照事件

CDC 可能产生两类不同语义:

  • 快照:读取某一时刻表中的当前行;
  • 增量变更:记录某次插入、更新或删除。

如果下游只按 updated_at 过滤,而没有考虑删除事件,就可能把已删除的数据永久保留在目标端。

如果更新事件只携带变更字段,则下游需要执行部分更新;如果携带更新后的完整行,则可以使用覆盖式 upsert。具体字段由 CDC 工具和配置决定,不能假设所有 CDC 事件格式相同。

4.2 顺序不是全局顺序

消息系统通常只能保证有限范围内的顺序,例如:

  • 同一分区内有序;
  • 同一 key 被路由到同一分区时,该 key 的事件有序;
  • 不同分区之间没有全局顺序。

因此,订单 o-10 的事件可能需要按版本号或日志位置判断:

order_id = o-10, version = 3, status = paid
order_id = o-10, version = 2, status = created

如果 version=3 先到,version=2 后到,简单的“最后到达覆盖”会把订单错误地恢复成 created

可靠的更新条件应类似:

UPDATE target_orders
SET status = :status,
    source_version = :version
WHERE order_id = :order_id
  AND source_version < :version;

这要求源事件拥有可比较的版本。时间戳可以作为辅助,但通常不如数据库日志位置、源端递增版本或事务序列可靠,因为时间戳可能相同、回拨或只表示应用写入时间。


五、幂等:重试后结果仍然正确

5.1 幂等的形式化定义

对操作 ff,如果满足:

f(f(x))=f(x)f(f(x)) = f(x)

则称 ff 对重复执行具有幂等性。

在数据管道中,xx 不仅是目标数据,还包括:

  • 当前目标状态;
  • 输入事件;
  • 事件身份;
  • 事件版本;
  • 转换规则。

例如,将事件 e-1001 写入订单事实表,重复写入不应产生两行订单记录。

但“没有重复行”并不自动等于幂等。以下两种操作可能产生不同结果:

INSERT INTO daily_sales(day, amount)
VALUES ('2025-01-01', 100);

重复执行后金额变成 200,不幂等。

INSERT INTO daily_sales(day, amount)
VALUES ('2025-01-01', 100)
ON CONFLICT (day)
DO UPDATE SET amount = EXCLUDED.amount;

重复执行后仍为 100,具备覆盖式幂等性。

5.2 事件身份、业务主键与版本

实现幂等通常需要区分三个概念:

  • 事件 ID:某条变更或消息的唯一身份;
  • 业务主键:目标实体的身份,例如 order_id
  • 版本:同一实体事件的先后关系。

例如:

event_id   = e-1001
order_id   = o-10
version    = 7

它们的作用不同:

  • event_id 防止同一事件重复处理;
  • order_id 决定写入哪条业务实体;
  • version 防止旧事件覆盖新状态。

只有业务主键而没有事件 ID,无法区分“重复事件”和“同一订单的新事件”;只有事件 ID 而没有版本,无法处理同一实体的乱序更新。

5.3 目标端唯一约束是重要防线

下面是一个 PostgreSQL 示例。目标表以 event_id 防止重复事件:

CREATE TABLE fact_order_events (
    event_id       text PRIMARY KEY,
    order_id       text NOT NULL,
    event_type     text NOT NULL,
    event_version  bigint NOT NULL,
    occurred_at    timestamptz NOT NULL,
    amount         numeric(18, 2),
    loaded_at      timestamptz NOT NULL DEFAULT now()
);

第一次写入:

INSERT INTO fact_order_events (
    event_id, order_id, event_type, event_version, occurred_at, amount
)
VALUES (
    'e-1001', 'o-10', 'paid', 7,
    '2025-01-01 10:00:00+00', 199.00
)
ON CONFLICT (event_id) DO NOTHING;

预期结果是插入一行。

再次执行完全相同的 SQL:

INSERT 0 0

目标表仍只有一行。ON CONFLICT 的成立依赖于 event_id 上存在唯一约束或主键。没有约束时,应用层的“先查询再插入”会在并发下产生竞态:

事务 A:查询不存在
事务 B:查询不存在
事务 A:插入
事务 B:插入

5.4 覆盖式 upsert 与追加式事实

对当前状态表,常用版本保护的 upsert:

CREATE TABLE current_orders (
    order_id       text PRIMARY KEY,
    status         text NOT NULL,
    source_version bigint NOT NULL,
    updated_at     timestamptz NOT NULL
);
INSERT INTO current_orders (
    order_id, status, source_version, updated_at
)
VALUES (
    'o-10', 'paid', 7, '2025-01-01 10:00:00+00'
)
ON CONFLICT (order_id) DO UPDATE
SET status = EXCLUDED.status,
    source_version = EXCLUDED.source_version,
    updated_at = EXCLUDED.updated_at
WHERE current_orders.source_version < EXCLUDED.source_version;

逐步看结果:

  1. 目标不存在:插入 version=7;
  2. 重复收到 version=7:WHERE 不成立,不更新;
  3. 收到 version=6:WHERE 不成立,旧事件不能覆盖新状态;
  4. 收到 version=8:WHERE 成立,更新为 version=8。

这适合“每个订单只保留当前状态”的表。

而对于审计明细或事件事实表,通常应追加每个合法事件,再通过 event_id 去重。不能把所有数据都做成当前状态,否则会丢失状态变化历史。

5.5 “恰好一次”需要明确边界

常说的 exactly-once 可能指不同事情:

  1. 消息系统中每条消息只交付一次;
  2. 计算引擎内部状态只更新一次;
  3. 目标数据库最终只产生一次业务效果;
  4. 对外部 API 的副作用只执行一次。

这些不是同一件事。

如果消费者写数据库和提交消息 offset 不在同一事务中,常见语义是:

  • at-least-once:可能重复消费;
  • 通过目标端幂等实现“重复执行但效果一次”。

真正的端到端 exactly-once 需要覆盖整个副作用边界。向不支持幂等键的外部 HTTP 服务发送扣款请求,即使消息系统内部支持事务,也不能自动保证扣款只发生一次。


六、一个可恢复的批处理设计

下面使用 PostgreSQL 展示“原始事件表 + 目标表 + 批次状态”的基本结构。示例不是完整调度器,但包含批处理最重要的事务边界。

6.1 输入和目标表

CREATE TABLE raw_order_events (
    event_id       text PRIMARY KEY,
    order_id       text NOT NULL,
    event_type     text NOT NULL,
    event_version  bigint NOT NULL,
    occurred_at    timestamptz NOT NULL,
    amount         numeric(18, 2),
    ingested_at    timestamptz NOT NULL DEFAULT now()
);

CREATE TABLE pipeline_runs (
    pipeline_name  text NOT NULL,
    batch_id       text NOT NULL,
    range_start    timestamptz NOT NULL,
    range_end      timestamptz NOT NULL,
    status         text NOT NULL CHECK (status IN ('running', 'succeeded', 'failed')),
    row_count      bigint,
    started_at     timestamptz NOT NULL DEFAULT now(),
    finished_at    timestamptz,
    error_message  text,
    PRIMARY KEY (pipeline_name, batch_id)
);

CREATE TABLE daily_order_sales (
    sales_day      date PRIMARY KEY,
    paid_amount    numeric(18, 2) NOT NULL,
    source_rows    bigint NOT NULL,
    pipeline_batch text NOT NULL,
    calculated_at   timestamptz NOT NULL DEFAULT now()
);

这里有三个层次:

  • raw_order_events:保留输入事件;
  • pipeline_runs:记录一次批处理的身份、范围和状态;
  • daily_order_sales:面向查询的汇总结果。

6.2 创建批次并处理

假设批次范围为:

[2025-01-01 00:00:00+00, 2025-01-02 00:00:00+00)

先登记批次:

INSERT INTO pipeline_runs (
    pipeline_name, batch_id, range_start, range_end, status
)
VALUES (
    'daily_order_sales',
    '2025-01-01-v1',
    '2025-01-01 00:00:00+00',
    '2025-01-02 00:00:00+00',
    'running'
)
ON CONFLICT (pipeline_name, batch_id) DO NOTHING;

然后在一个事务中重新计算这个批次的结果:

BEGIN;

DELETE FROM daily_order_sales
WHERE sales_day >= DATE '2025-01-01'
  AND sales_day <  DATE '2025-01-02';

INSERT INTO daily_order_sales (
    sales_day,
    paid_amount,
    source_rows,
    pipeline_batch
)
SELECT
    occurred_at::date AS sales_day,
    COALESCE(SUM(amount), 0) AS paid_amount,
    COUNT(*) AS source_rows,
    '2025-01-01-v1'
FROM raw_order_events
WHERE occurred_at >= '2025-01-01 00:00:00+00'
  AND occurred_at <  '2025-01-02 00:00:00+00'
  AND event_type = 'paid'
GROUP BY occurred_at::date;

UPDATE pipeline_runs
SET status = 'succeeded',
    row_count = (
        SELECT COUNT(*)
        FROM raw_order_events
        WHERE occurred_at >= '2025-01-01 00:00:00+00'
          AND occurred_at <  '2025-01-02 00:00:00+00'
          AND event_type = 'paid'
    ),
    finished_at = now()
WHERE pipeline_name = 'daily_order_sales'
  AND batch_id = '2025-01-01-v1'
  AND status = 'running';

COMMIT;

这个方案的核心是分区或时间范围替换

  • 先删除本批次负责的目标范围;
  • 再从原始层完整计算;
  • 计算成功后才提交;
  • 如果中途失败,事务回滚,旧结果仍然存在。

因此,批次重跑不会把金额重复累加。

6.3 这个方案的边界

上述方案成立需要满足:

  1. raw_order_events 在重算期间可读;
  2. 批次范围定义稳定;
  3. 目标表中的该范围可以被完整替换;
  4. 下游不会在事务提交前读取到半成品;
  5. 迟到数据的处理策略明确。

如果目标是分布式对象存储上的多个文件,无法依赖 PostgreSQL 事务保护全部副作用,就需要使用:

  • 临时目录写入;
  • 文件清单或 manifest;
  • 原子提交标记;
  • 版本化表;
  • 分区交换;
  • 失败清理和孤儿文件扫描。

“删除目标分区再写入”并不天然安全。如果删除已提交、写入尚未完成,任务失败时就会留下空分区。必须让删除和替换处于同一事务,或采用版本化提交策略。


七、流处理中的状态与失败路径

流处理通常可抽象为:

读取消息
  -> 解析
  -> 校验
  -> 去重 / 版本判断
  -> 转换
  -> 写入目标
  -> 提交消费位置

每一步都可能失败。

7.1 解析失败

例如消息不是合法 JSON,或者字段类型错误。常见做法是:

  • 不断重试同一条消息;
  • 将消息写入死信队列;
  • 记录原始 payload、错误原因和消息位置;
  • 继续处理后续消息。

如果解析错误消息不可能通过重试恢复,持续重试会阻塞整个分区。死信记录必须包含足够信息,使其可以独立重放。

7.2 目标写入成功、offset 提交失败

这是最常见的重复路径:

读取 e-1001
  -> 写入目标成功
  -> offset 提交失败
  -> 重启后再次读取 e-1001
  -> 目标端幂等写入

目标端必须有唯一约束或去重表。可以额外保存处理记录:

CREATE TABLE processed_events (
    event_id       text PRIMARY KEY,
    processed_at   timestamptz NOT NULL DEFAULT now(),
    payload_hash   text NOT NULL
);

处理时将“登记事件”和“业务写入”放入同一数据库事务:

BEGIN;

INSERT INTO processed_events(event_id, payload_hash)
VALUES ('e-1001', 'sha256:abc')
ON CONFLICT (event_id) DO NOTHING;

应用需要检查插入影响行数:

  • 插入 1 行:首次处理,继续写业务表;
  • 插入 0 行:事件已处理,跳过业务写入。

如果同一个 event_id 再次出现但 payload_hash 不同,应视为数据冲突,而不是普通重复。它可能表示源端重用了事件 ID、序列化不一致或数据被篡改。

7.3 目标写入失败、offset 已提交

如果先提交 offset、后写目标,目标写入失败后消息可能不会再次出现,造成数据丢失。因此通常应遵循:

目标结果成功持久化
        |
        v
再提交消费位置

但如果目标是外部系统,提交 offset 与外部副作用无法组成原子事务,就需要使用目标端幂等键、事务性 outbox、重试表或人工对账。


八、补数:重新处理历史数据,而不是简单重跑任务

补数是对既有时间范围、分区、主键集合或事件集合进行重新计算和写入。补数原因包括:

  • 上游漏发或迟到;
  • CDC 消费中断;
  • 转换逻辑有缺陷;
  • 维表修正;
  • 业务规则变更;
  • 目标分区损坏;
  • 历史数据回溯。

8.1 补数范围必须可描述

一个补数任务至少应记录:

pipeline: daily_order_sales
range: [2025-01-01, 2025-01-08)
reason: 修复税额转换错误
code_version: v2.3.1
source_snapshot: 2025-01-08T02:00Z

没有明确范围的“全量重跑”风险很高:

  • 可能覆盖仍在生产中的新结果;
  • 可能使用了与历史不同的维表版本;
  • 可能扩大锁和资源影响;
  • 可能让问题难以回溯。

8.2 重算还是增量修补

有两种主要方式。

重算目标范围

对目标日期或分区执行:

删除或替换目标范围
从原始层重新聚合
校验
提交新版本

优点是逻辑简单,天然避免累加重复。

增量修补

只计算差异,例如:

旧结果:1000
补入事件:+50
撤销事件:-20
新结果:1030

增量修补要求正确维护:

  • 原事件是否已经计入;
  • 修补是否执行过;
  • 撤销与反向事件;
  • 重复补数;
  • 并发更新。

如果差异识别不可靠,增量修补会把错误放大。对可重算的聚合,通常优先选择范围重算;只有在数据量或延迟要求使重算不可接受时,才采用增量修补。

8.3 迟到事件与补数窗口

流系统常使用“滚动补数窗口”,例如每次重算最近 7 天:

窗口 = [当前时间 - 7 天, 当前时间)

这样可以吸收迟到事件,但不能解决无限迟到。窗口大小应根据:

  • 源系统最大延迟;
  • CDC 中断恢复时间;
  • 业务允许的修正周期;
  • 下游报表的锁定时间。

财务结算日后仍允许调整时,窗口不能简单按自然日关闭;应引入“结算状态”和显式更正流程。


九、数据校验:证明处理结果可信

校验不是只执行一次 COUNT(*)。完整校验应覆盖:

  1. 完整性;
  2. 唯一性;
  3. 合法性;
  4. 一致性;
  5. 汇总可对账性;
  6. 时效性。

9.1 行数校验

最基本的校验:

SELECT COUNT(*) AS source_rows
FROM raw_order_events
WHERE occurred_at >= '2025-01-01 00:00:00+00'
  AND occurred_at <  '2025-01-02 00:00:00+00'
  AND event_type = 'paid';

目标侧:

SELECT SUM(source_rows) AS target_source_rows
FROM daily_order_sales
WHERE sales_day >= DATE '2025-01-01'
  AND sales_day <  DATE '2025-01-02';

但两者不一定直接相等:

  • 源数据按事件行统计;
  • 目标可能按订单去重;
  • 一个订单可能有多条支付事件;
  • 目标可能过滤无效记录。

所以校验必须比较同一语义。如果目标按订单去重,源侧也应:

SELECT COUNT(DISTINCT order_id)
FROM raw_order_events
WHERE ...;

9.2 金额和边界校验

金额聚合可以做:

SELECT
    COUNT(*) AS row_count,
    COALESCE(SUM(amount), 0) AS amount_sum,
    MIN(amount) AS min_amount,
    MAX(amount) AS max_amount
FROM raw_order_events
WHERE occurred_at >= '2025-01-01 00:00:00+00'
  AND occurred_at <  '2025-01-02 00:00:00+00'
  AND event_type = 'paid';

需要注意 NULL

  • COUNT(*) 统计所有行;
  • COUNT(amount) 不统计 amount IS NULL
  • SUM(amount) 忽略 NULL;
  • 没有行时 SUM 可能返回 NULL,因此示例使用 COALESCE

金额校验还需要明确精度。数据库中的 numericDECIMAL 与浮点数的加法语义不同。财务金额不应在转换过程中随意使用二进制浮点数后再比较精确相等。

9.3 校验和

可以构造行级规范化表示,再计算哈希。关键是必须定义规范化规则:

  • 字段顺序;
  • NULL 表示;
  • 字符编码;
  • 时间时区;
  • 数字格式;
  • 换行与空格;
  • 是否包含被忽略字段。

例如,对以下逻辑行:

order_id=o-10|status=paid|amount=199.00

计算哈希后,再对分区内哈希做汇总。

但普通哈希的 XOR 或简单求和存在风险:

  • XOR 对重复行可能抵消;
  • 求和可能发生碰撞或溢出;
  • 字段格式不一致会造成误报;
  • 哈希相同也不是数学上的绝对证明。

因此,实践中通常组合使用:

行数 + 主键范围 + 金额汇总 + 分区哈希 + 抽样明细

校验和适合发现差异,不应被当作绝对证明。

9.4 唯一性和引用一致性

目标表有主键时,数据库会阻止重复主键;没有约束的文件或宽表则要主动检查:

SELECT order_id, COUNT(*)
FROM current_orders
GROUP BY order_id
HAVING COUNT(*) > 1;

外键关系也需要检查。例如订单明细不能引用不存在的订单:

SELECT COUNT(*) AS orphan_rows
FROM order_items i
LEFT JOIN current_orders o ON o.order_id = i.order_id
WHERE o.order_id IS NULL;

如果目标数据库没有强制外键,数据管道必须通过校验任务承担这部分责任。

9.5 时间性校验

一批数据可能“数量正确但已经过时”。应记录:

  • 最早事件时间;
  • 最晚事件时间;
  • 最大摄入延迟;
  • 当前 watermark;
  • 源端日志位点;
  • 目标端最新成功时间。

例如:

latency=target_loaded_atevent_occurred_at\text{latency} = \text{target\_loaded\_at} - \text{event\_occurred\_at}

延迟的平均值可能正常,但 P99 严重恶化,因此通常要观察分位数,而不是只看平均值。


十、可观测性:让数据问题能够被解释

数据管道的可观测性至少包含三类信号:

  • 指标;
  • 日志;
  • 追踪和关联信息。

10.1 指标

应按管道、批次、分区和数据源记录:

吞吐

input_rows
output_rows
input_bytes
output_bytes
messages_consumed

延迟

processing_latency
source_to_target_latency
consumer_lag
watermark_lag

错误

parse_error_count
validation_error_count
dead_letter_count
retry_count
upsert_conflict_count

数据质量

null_rate
duplicate_rate
orphan_rate
amount_delta
row_count_delta
checksum_mismatch

指标必须带有限基数的标签,例如 pipelinesourcepartition。不要把 event_id 作为监控指标标签,否则高基数会使监控系统本身失控。单个事件 ID 应放在日志或追踪上下文中。

10.2 结构化日志

一条可用的失败日志至少包含:

{
  "pipeline": "daily_order_sales",
  "run_id": "2025-01-01-v1",
  "source": "raw_order_events",
  "range_start": "2025-01-01T00:00:00Z",
  "range_end": "2025-01-02T00:00:00Z",
  "partition": "p-03",
  "event_id": "e-1001",
  "order_id": "o-10",
  "source_position": "123456",
  "code_version": "v2.3.1",
  "error_type": "validation_error",
  "error": "amount must not be null"
}

其中 run_idbatch_idevent_id、CDC 位点和代码版本是定位重跑边界的关键。

10.3 数据血缘

血缘回答:

这个指标来自哪些表、哪些字段、哪些批次和哪些规则?

至少应能关联:

源表 / 消息
  -> 原始层批次
  -> 转换任务版本
  -> 目标表分区
  -> 下游报表或指标

如果发现目标收入异常,不能只知道 SQL 任务失败;还要知道:

  • 哪些源分区参与了计算;
  • 使用了哪个规则版本;
  • 哪些事件被过滤;
  • 是否发生过补数;
  • 下游是否已读取旧版本结果。

10.4 告警条件

好的告警应对应可执行动作。例如:

  • 消费 lag 持续超过阈值:检查消费者、分区热点和目标写入;
  • 死信数量增加:检查 schema 演进或脏数据;
  • 行数差异超阈值:暂停下游发布并启动对账;
  • 批次超过 SLA:检查锁、资源、外部依赖;
  • watermark 长时间不推进:检查乱序数据或分区阻塞;
  • checksum 不匹配:定位具体分区和批次,不要直接全量重跑。

只告警“任务失败”通常不足以诊断。任务成功但数据少了 30% 更危险,因此数据质量告警必须独立于任务运行状态。


十一、Schema 演进与契约

数据管道依赖字段结构。字段变化包括:

  • 新增字段;
  • 字段重命名;
  • 类型扩大或缩小;
  • 枚举值增加;
  • 字段变为 NULL;
  • 删除字段;
  • 嵌套结构变化。

新增可选字段通常较容易兼容;删除字段和类型缩小更危险。

例如源端将:

amount: integer

改为:

amount: decimal(18,2)

如果下游仍按整数解析,可能出现:

  • 解析失败;
  • 小数被截断;
  • 金额精度错误。

CDC 系统还涉及 schema history、消息格式和消费者兼容策略。不能仅依赖“数据库 DDL 执行成功”判断下游兼容。

一种可审计的处理方式是为输入记录 schema 版本:

event_id
schema_version
payload

转换任务明确支持哪些版本,并对不支持的版本进入隔离区,而不是静默丢弃字段。


十二、事务边界:源端、管道和目标端不是一个事务

数据库事务只保护它所在的事务资源。

12.1 单数据库内

如果读取原始表、写入目标表和更新批次状态都在同一个 PostgreSQL 数据库中,可以使用一个事务保证:

目标结果提交
批次状态成功提交

两者一起成功或一起回滚。

PostgreSQL 的 INSERT ... ON CONFLICT、事务隔离和约束可以帮助实现这种原子性。MySQL 8.4 也支持事务、唯一约束以及 INSERT ... ON DUPLICATE KEY UPDATE 等机制,但 SQL 语法和具体冲突行为不同,不能直接把 PostgreSQL SQL 原样搬过去。

12.2 跨数据库或跨系统

如果源端是 MySQL、目标端是 PostgreSQL,中间还有消息队列,则通常不存在一个自动覆盖所有系统的本地事务:

MySQL 事务
   != 消息提交
   != PostgreSQL 事务
   != 外部 API 调用

此时需要设计:

  • CDC 位点记录;
  • 目标端幂等;
  • 事务性 outbox;
  • 可重放原始层;
  • 对账和补偿;
  • 明确的失败状态。

分布式事务并不能消除业务语义问题。即使技术上能协调多个资源,也仍需决定迟到事件、版本冲突和错误数据如何处理。


十三、常见错误与诊断路径

错误一:用 updated_at 作为唯一增量边界

问题:

WHERE updated_at > :last_time

如果多个事务产生相同时间戳,或者时钟精度不足,可能遗漏数据。若任务在读取期间发生更新,也可能出现边界不稳定。

改进方式:

  • 使用稳定且可比较的日志位点;
  • 使用 (updated_at, primary_key) 复合游标;
  • 对时间边界保留重叠窗口;
  • 通过幂等 upsert 消化重叠数据。

错误二:先更新 offset,再写目标

表现:

  • 消费位置持续推进;
  • 目标行数低于源端;
  • 重启后无法自动重新获取已提交消息。

诊断:

  • 对比消息 offset 与目标批次;
  • 查看目标写入失败日志;
  • 检查 offset 提交时间是否早于目标提交时间。

恢复:

  • 从原始消息、CDC 日志或源库重新建立读取范围;
  • 不要只依赖当前 offset;
  • 重新处理时必须依赖目标幂等。

错误三:用“最后到达”处理更新

表现:

目标状态从 paid 回退为 created

原因通常是乱序消息覆盖。

诊断:

SELECT order_id, source_version, status
FROM current_orders
WHERE order_id = 'o-10';

再对比该订单所有输入事件的版本顺序。恢复时应按源版本重放,或执行带版本保护的更新。

错误四:聚合任务使用累加写入

错误模式:

UPDATE daily_sales
SET amount = amount + :batch_amount
WHERE sales_day = :day;

任务重试会再次累加。

修复方式:

  • 按范围删除后重算;
  • 使用批次版本覆盖;
  • 为每个输入事件去重;
  • 保存增量应用记录并确保其与聚合更新同事务提交。

错误五:把任务成功当作数据正确

任务可能成功完成,但输入本身已缺失,或者过滤条件发生变化。

必须区分:

运行成功:程序没有异常退出
数据成功:结果通过完整性和业务校验

只有第二种状态才适合发布给下游。


十四、一个实用的端到端状态模型

可以为每个批次或事件维护如下状态:

received
  -> validated
  -> transformed
  -> written
  -> reconciled
  -> published

失败状态则包括:

rejected
retrying
dead_letter
failed
compensating

例如一个批次的状态转移:

  1. received:已发现输入范围;
  2. validated:schema、范围和基本质量检查通过;
  3. transformed:中间结果生成;
  4. written:目标数据提交成功;
  5. reconciled:行数、金额和主键检查通过;
  6. published:下游可以读取该版本。

只有记录到 published,下游才应认为数据可用。否则,“目标表中已经存在部分行”不应等于“批次已完成”。


十五、如何选择处理策略

可以按数据特征选择:

使用批处理重算,适合:

  • 数据天然按日、小时或分区组织;
  • 目标结果可以按范围替换;
  • 允许分钟级或小时级延迟;
  • 规则变更后需要重放;
  • 聚合逻辑复杂但可确定性重算。

使用流处理,适合:

  • 需要低延迟更新;
  • 输入持续到达;
  • 业务事件有明确身份和版本;
  • 可以接受增量状态维护;
  • 已设计迟到、乱序和重放机制。

使用批流混合,适合:

  • 流处理提供实时近似结果;
  • 批处理定期校正历史结果;
  • 下游同时需要实时看板和稳定报表;
  • 事件量大,无法频繁全量重算。

无论选择哪种方式,都应让“处理结果可重放”。实时路径追求延迟,批路径负责校正,两者必须有可比较的结果和校验规则。


十六、最终检查标准

一条数据管道达到可生产运行的基本条件,不是“调度成功”,而是能够回答:

  • 输入范围是否明确且可重放?
  • 每个事件或批次是否有稳定身份?
  • 重复消费会不会改变最终结果?
  • 乱序事件是否按版本处理?
  • 目标写入与进度提交的失败路径是什么?
  • 迟到数据如何影响窗口和报表?
  • 补数是否有范围、原因、代码版本和审批记录?
  • 结果是否通过行数、金额、唯一性和时效性校验?
  • 错误数据是否进入可查询、可重放的隔离区?
  • 指标、日志和血缘是否能定位到具体批次或事件?
  • Schema 变化是否有兼容策略?
  • 下游读取的是已校验并发布的版本,还是半成品?

ETL 与 ELT 的区别决定数据在哪里转换,批处理与流处理决定数据如何推进,幂等决定失败重试是否安全,补数决定历史错误能否修复,校验决定结果是否可信,可观测性则决定问题能否被发现和解释。只有这些机制共同成立,数据管道才不仅能“跑起来”,还能在重复、失败、迟到、乱序和规则变化中保持正确。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。