千万级数据同步与 ETL 实战(二):快照、Watermark 与无漏数分页
增量查询写成 WHERE update_time > ? 很容易,难点集中在三个交界处:首次快照进行时源表仍在写;同一时间戳容纳多行;目标端已提交而 Checkpoint 尚未提交。任一处处理不清,系统都会在低负载测试中正常运行,在真实并发和故障下出现漏数或旧值覆盖。
本文从一个约束明确的例子展开。源表有 1000 万行,主键 id,更新时间精确到毫秒,业务允许 1 分钟调度周期。源端只提供普通查询权限,没有日志消费能力。目标端支持按业务键 Upsert。
一、首次快照需要 Epoch
快照不是简单的 SELECT *。它是目标数据集的一次基线构建,必须能和快照期间发生的变化衔接。为每次基线分配单调递增的 snapshotEpoch:
task = customer_sync
snapshotEpoch = 42
schemaId = customer-v7Epoch 42 的每个 Batch 都引用 customer-v7,其编码口径在本轮基线构建期间保持冻结。新的快照开始后,旧 Epoch 仍可继续服务,直到新 Epoch 完成校验并原子切换。直接清空目标表再写入,会把数小时的快照过程暴露给查询方,也失去快速回退能力。
完整流程如下:
如果数据库支持一致性快照和可对应的日志位置,H0 可以是事务快照对应的 LSN、SCN 或 Binlog Position。快照读取看到边界时刻的一致状态,增量从该位置之后消费,衔接关系能够精确定义。
只有 update_time 时,情况更弱。捕获 H0=currentTime 后读取全表,不能证明所有 update_time <= H0 的行都处于同一个数据库快照。长快照期间的更新、事务提交延迟和时钟偏差会让边界出现重叠或空隙。可行做法是保留重叠区间,从 H0-overlap 开始反复读取,快照结束后追平,再执行一次 Warm 或 Cold 对账。文章和系统界面都应标注这是收敛性补偿策略,不应宣称取得了日志级无缝切换。
二、Watermark 必须包含稳定键
只保存 lastUpdateTime 会在同时间戳分页时出错。假设 4000 行记录的更新时间均为 10:00:00.123,页大小为 1000。读取第一页后任务中断,如果 Checkpoint 保存该时间,恢复查询使用 update_time > lastUpdateTime,剩余 3000 行永远不会再被读取。
Checkpoint 应保存二元游标:
$$ W=(t,k) $$
其中 $t$ 为更新时间,$k$ 为稳定且唯一的主键。比较采用字典序:
$$ (t_1,k_1)<(t_2,k_2) \iff t_1<t_2\ \lor\ (t_1=t_2\land k_1<k_2) $$
对应 SQL:
SELECT id, name, status, update_time
FROM source_table
WHERE update_time > :last_time
OR (update_time = :last_time AND id > :last_id)
ORDER BY update_time, id
LIMIT :page_size;(update_time, id) 需要组合索引。没有合适索引时,数据库可能每轮排序或扫描大量历史行,应用层页大小调得再小也不会降低源端工作量。
复合主键需要完整的字典序条件。以 (tenant_id, order_id) 为例,不能只保存 order_id:
WHERE update_time > :t
OR (
update_time = :t
AND (
tenant_id > :tenant
OR (tenant_id = :tenant AND order_id > :order_id)
)
)
ORDER BY update_time, tenant_id, order_id
LIMIT :page_size;数据库支持行值比较时可以写成 (update_time, tenant_id, order_id) > (?, ?, ?),但要核对数据库版本、索引使用情况和 NULL 语义。参与游标的列应为非空,排序规则也要固定。
三、分页为什么不能使用 OFFSET
LIMIT 1000 OFFSET 9000000 需要数据库跳过前 900 万行,深分页成本随偏移增长。并发修改还会改变记录位置。第一页读取后新增一行排在前面,第二页的偏移可能重复一条旧记录;删除一行则可能跳过原本处于页首的数据。
Keyset Pagination 使用上一页末尾键定位下一页:
SELECT id, payload
FROM source_table
WHERE id > :last_id
ORDER BY id
LIMIT :page_size;索引定位成本约为 $O(\log N)$,随后顺序读取 $P$ 条,单页可近似写成:
$$ T_{page}=O(\log N + P) $$
如果主键是随机 UUID,Keyset 仍能保证稳定顺序,但局部性和索引页读取成本可能较差。可以按 UUID 前缀建立逻辑分片,每个分片内部继续使用完整 UUID 游标。若源数据没有任何稳定唯一键,就无法保证重启后从同一位置继续,需要先配置业务键或接受全量重读和结果去重。
四、重叠窗口处理迟到提交
严格从已提交 Watermark 之后读取,仍会漏掉“业务记录晚到、时间值偏旧”的数据。常见来源包括事务持续时间、节点时钟偏差、时间精度截断、人工回填历史记录,以及应用先写业务字段后补写更新时间。
每轮查询可使用两个边界:
$$ low=W_{committed}.time-overlap $$
$$ high=now-safetyDelay $$
例如调度周期 60 秒、重叠 5 分钟、安全延迟 30 秒。系统每轮回读已提交边界前 5 分钟,暂不读取离当前时间 30 秒以内的数据。重叠带来重复,目标端必须按身份键和源版本幂等处理。
overlap 不能凭感觉设置。可以统计源端提交时间与业务时间的差值 $L=t_{visible}-t_{business}$,选择覆盖目标分位数的窗口:
$$ overlap \ge Q_{0.999}(L)+clockSkew+precisionError $$
窗口只能覆盖有限迟到。人工把一年前的记录改完仍保留原时间,5 分钟重叠无法发现,必须由 Hash Wheel 或全量校验兜底。
五、读完一页不等于可以推进 Watermark
抽取侧读到一页数据并不代表同步完成。批次可能尚未落盘,目标端也可能写入失败。安全顺序应为:
读取 Source Page
→ 生成不可变 Batch 和 Manifest
→ Batch 持久化完成
→ 目标端幂等提交
→ 校验提交结果
→ CAS 更新 Checkpoint如果 Batch 还未持久化就推进 Watermark,进程宕机后该页无法重放。若目标端提交成功、Checkpoint 更新失败,恢复后会再次读取并写入同一批数据。这条路径依赖幂等 Sink,结果可以保持正确。
Checkpoint 表需要版本号,以比较并交换方式阻止两个 Worker 同时推进:
UPDATE etl_checkpoint
SET cursor_value = :next_cursor,
last_batch_id = :batch_id,
version = version + 1,
committed_at = CURRENT_TIMESTAMP
WHERE task_id = :task_id
AND partition_id = :partition_id
AND version = :expected_version;影响行数为 0 表示租约失效或并发冲突。当前 Worker 不能覆盖新 Checkpoint,应停止分片并重新加载状态。
六、Source Version 决定旧值能否覆盖新值
每条变化应携带可比较的源版本。单调业务版本优先,日志位置其次。更新时间可能相同,因此可组合主键形成稳定序列:
sourceVersion = encode(update_time, primary_key)这只保证同一排序流内的顺序。如果记录 A 在 10:05 被更新,随后在 10:06 又被人工写入旧的 update_time=09:00,版本比较会拒绝这次业务上较新的修改。此类源系统需要数据库版本、触发器维护的行版本,或由 Warm 通道发现差异后按权威源覆盖。设计文档要把时间字段是否单调作为数据契约,不能默认成立。
目标表可保存当前版本和来源 Epoch:
INSERT INTO target_table(id, payload, source_version, epoch)
VALUES (:id, :payload, :version, :epoch)
ON CONFLICT (id) DO UPDATE
SET payload = EXCLUDED.payload,
source_version = EXCLUDED.source_version,
epoch = EXCLUDED.epoch
WHERE target_table.source_version <= EXCLUDED.source_version;不同数据库的 Upsert 语法和并发语义不同,需要在目标数据库上做并发测试,不能只检查 SQL 能执行。
七、Append-Only 只能覆盖新增
只有自增主键、且历史记录严格不变时,可以保存 lastId:
SELECT id, payload
FROM source_table
WHERE id > :last_id
ORDER BY id
LIMIT :page_size;读取量接近新增量 $\Delta$。这个模式无法发现旧行修改和删除。任务配置要显式声明 appendOnly=true,并由数据所有者确认约束。若约束无法保证,可以保留高频 lastId 通道处理新增,再用低频 Hash Wheel 检查历史数据。
自增 ID 也可能因批量导入、主从切换或序列重置而回退。监控应记录源端 max(id),发现小于已提交 lastId 时暂停任务,避免在错误假设下继续运行。
八、HTTP 与文件的游标语义
HTTP 数据源优先使用服务端返回的 nextCursor、pageToken 或版本字段。只有页码时,并发新增会造成页漂移,风险与数据库 OFFSET 相同。ETag 和 Last-Modified 适合判断资源是否变化,但弱 ETag、缓存节点和服务端实现会影响语义,不能直接当作记录级版本。
文件数据源的 Checkpoint 可以由 (path, versionId, size, mtime, contentHash) 组成。对象存储的 Multipart ETag 未必等于文件 MD5。需要强内容校验时,应读取平台提供的校验字段或自行计算摘要。大文件可以先比较文件版本,再比较 Block Hash,随后定位到记录级差异。
无论游标来自数据库、HTTP 还是文件,Checkpoint 都应保存类型、编码版本和任务 Epoch。把所有游标塞进一段无版本字符串,后续变更排序规则或时区时很难迁移。
九、用故障场景验收衔接逻辑
对 1000 万行合成数据运行首次快照,并在快照期间持续新增、更新和删除。测试至少注入这些故障:第三页读取后断网、Batch 落盘一半宕机、目标端提交后进程退出、Checkpoint CAS 冲突、源端时钟回拨、数千行使用同一毫秒时间戳。
验收不能只比较汇总行数。需要比较每个业务键的值、版本、删除状态和 Schema,确认:
$$ Target_{final}=Expected(Source, boundary, deletePolicy) $$
其中 boundary 和 deletePolicy 必须写入测试定义。物理删除若只能由每小时一次的 Warm 通道发现,紧接快照完成时比较会得到预期内的短暂差异。只有把时间边界和保证范围写进断言,测试结果才可解释。