千万级数据同步与 ETL 实战(一):边界、降级策略与整体架构
数据同步项目开工时,团队常先讨论 Kafka、Flink 或 Debezium。这个顺序容易忽略一个基础问题:源系统究竟能提供什么。很多企业数据源只有只读账号、分页接口或定期投递的文件。数据库日志权限拿不到,表里也未必有可靠的更新时间。此时直接套用 CDC 架构,实施阶段会在权限、版本、网络和运维能力上不断返工。
本文先给系统划定边界。目标数据集以单表或单集合不超过 1000 万行为基准,变更速率不高,业务接受几十秒到数小时的延迟。系统需要发现新增、修改和删除,支持断点恢复,并限制对源库的影响。每秒数十万事件、复杂事件时间计算、强实时交易链路不在这套微批架构的处理范围内,第六篇会单独讨论真正的流式系统。
一、先盘点数据源能力
同一个“查询数据”接口,能提供的增量能力差别很大。建任务前要探测并固化以下信息:
| 能力 | 需要确认的内容 | 缺失后的影响 |
|---|---|---|
| 变更日志 | Binlog、WAL、LSN、SCN、日志保留时间 | 无法按日志位置持续消费 |
| 版本字段 | 单调版本号、精度、是否覆盖删除 | 只能依赖时间或扫描 |
| 更新时间 | 是否有索引、是否随每次修改更新 | 历史回改可能漏读 |
| 稳定键 | 单主键、复合业务键、键是否会变化 | 无法可靠识别同一条记录 |
| 一致性快照 | 隔离级别、快照时长、日志衔接点 | 首次全量与增量可能断层 |
| 分页方式 | 游标、Keyset、页码、排序稳定性 | 并发写入时可能漏读或重复 |
| 删除信息 | 删除日志、软删字段、墓碑记录 | 需要周期扫描才能识别删除 |
| 源端改造 | 能否加索引、生成列、触发器 | 分桶扫描可能退化成全表计算 |
能力探测结果进入数据集版本,不能只留在部署人员的经验里。任务发布后如果主键、时区或更新时间精度发生变化,旧 Checkpoint 可能失效。调度器需要暂停任务,完成迁移验证后再恢复。
建议把增量策略设计成一条降级链:
降级需要记录原因。例如数据库支持 Binlog,但账号缺少复制权限,任务可以暂时使用 Watermark,同时持续告警“当前链路无法直接识别物理删除”。若系统只把模式保存为一个枚举,维护人员看不到能力缺口,也无法判断将权限补齐后能获得什么收益。
二、先承认变化检测的下界
设源表有 $N$ 行,没有日志、版本号、更新时间、触发器和预先维护的摘要。算法跳过了记录 $r$,只读取其余 $N-1$ 行。构造两个数据集:
$$ D_0 = {r, x_2, \ldots, x_N} $$
$$ D_1 = {r’, x_2, \ldots, x_N}, \quad r \ne r' $$
算法在两个数据集上得到的读取结果完全相同,因此无法判断 $r$ 是否变化。要保证发现任意一行变化,就要检查每个可能变化的位置,完整检测下界为:
$$ T(N)=\Omega(N) $$
Hash Wheel、Merkle Tree 和 IBLT 都不能消除这个下界。它们优化的是扫描如何分摊、摘要如何合并、差异如何定位、网络传多少数据,以及失败后重算多大范围。写架构方案时如果宣称“无日志、无时间字段也能只扫描变化行”,该结论缺少信息来源。
这个下界还能推导一个容量约束。设完整扫描周期的目标为 $C$ 秒,源端允许的持续扫描能力为 $R$ 行每秒,则必须满足:
$$ R \times C \ge N $$
1000 万行若要求 1 小时内发现任意历史修改,持续读取能力至少约为 2778 行每秒,还未计入重试、热点倾斜和目标端阻塞。若源库只能稳定提供每秒 500 行,完整周期下限接近 5.6 小时。此时要调整延迟目标、增加源端索引,或取得 CDC 能力,不能靠调度参数绕过吞吐约束。
三、控制面与数据面分开
任务配置、版本发布、分片租约和 Checkpoint 属于控制面;读取、规范化、比较、转换和写入属于数据面。两者放在一个执行循环里,常见后果是 Worker 重启后既不知道应该执行哪个版本,也不知道哪个批次已经对外生效。
控制面发布的是不可变任务版本,至少包括数据源能力快照、字段投影、身份键、规范化规则、增量模式、分片算法、目标写入语义和 Schema 版本。Worker 领取分片时同时取得任务版本和 Epoch。执行中的任务不读取一份会被随时覆盖的 JSON 配置,否则同一批数据可能前半段按旧规则计算指纹,后半段按新规则计算。
分片数量可以多于 Worker 数量。例如 128 个逻辑分片由 8 个 Worker 通过租约领取。逻辑分片决定状态和恢复粒度,Worker 数量决定当前并发。扩容时无需重建分片;单个分片失败时,其余分片仍可推进。
四、统一变化记录的边界
数据库行、HTTP 对象和文件记录统一转换成变化模型。模型要能回答五个问题:数据属于哪个任务版本、来自哪个分区、同一条数据如何识别、它在源端的版本是什么、它属于哪个可重放批次。
public record ChangeRecord(
String taskId,
long epoch,
String dataset,
String partitionId,
String identityKey,
Operation operation,
String sourceVersion,
String schemaId,
byte[] fingerprint,
Instant sourceTime,
Instant extractedAt,
Map<String, Object> before,
Map<String, Object> after,
String batchId,
long sequence
) {}
public enum Operation {
SNAPSHOT, INSERT, UPDATE, DELETE
}identityKey 和 fingerprint 不能合并。前者回答“是否为同一条记录”,后者回答“内容是否改变”。如果把整行摘要当作身份,修改一列会被解释成删除旧记录再新增记录,外键关系、审计链和下游聚合都会受到影响。
sourceVersion 也不能默认使用抽取时间。抽取时间只表示 Worker 何时读到记录,无法阻止迟到的旧值覆盖新值。这个字段必须来自源端可验证的先后关系;第二篇会结合 Watermark 讨论具体编码。只有行指纹可用时,目标端具备幂等性,却未必能比较新旧顺序,需要由同一分区串行处理,或在冲突时回源核对。
五、三条校验通道
只有一条 Watermark 通道时,人工回写旧时间、软删字段遗漏更新、应用时钟回拨都会形成长期差异。可以按不同成本设置三条通道:
Hot 通道承担日常低延迟同步,读取成本接近变化量 $\Delta$。Warm 通道按桶扫描,确保每个稳定键在规定周期内被重新检查。Cold 通道低频重算全量摘要,负责发现状态库损坏、分桶算法缺陷和长期漂移。三条通道输出相同的 ChangeRecord,共享幂等键,但拥有独立 Checkpoint,避免一条补偿链阻塞日常增量。
重复发现同一变化并非错误。假设 Hot 通道已写入版本 120,Warm 通道稍后又产生版本 120,目标端按 (taskId, identityKey, sourceVersion) 去重即可。真正危险的是通道之间没有版本比较,导致旧扫描结果覆盖已提交的新值。
六、容量与延迟预算
以 1000 万行、平均原始记录 1KB 为例,源数据约 10GB。指纹状态若保存 16 字节身份摘要、16 字节内容指纹、8 字节 Generation 和 16 至 32 字节版本元数据,裸数据约 560MB 至 720MB。KV 引擎还会产生索引、日志和 Compaction 开销,第三篇会给出完整磁盘估算。
这个数字只用于容量起点。主键长度、写放大、压缩率和保留的旧版本都会改变结果。上线前应生成接近真实键分布的合成数据,持续跑完至少两个 Compaction 周期,再决定磁盘水位和扩容阈值。
延迟也要分开测量:
$$ L_{end}=L_{detect}+L_{queue}+L_{transform}+L_{sink}+L_{commit} $$
Watermark 周期设为 60 秒,只约束变化发现阶段。目标端持续阻塞时,队列等待会主导总延迟。监控只展示调度是否准时,会掩盖待处理批次已经积压数小时的事实。
七、明确不保证的内容
微批系统可以通过不可变 Batch、幂等写入和延迟提交 Checkpoint 获得可恢复性,但不能自动获得跨源库、状态库和目标库的分布式事务。源数据在长时间快照期间持续更新时,如果源端没有一致性快照或日志衔接点,首次全量和增量之间仍要靠重叠读取与全量对账收敛。
删除检测也有延迟边界。Watermark 通常发现不了物理删除,只有日志、软删标记或完整桶扫描能给出证据。桶扫描未完成时禁止生成删除事件,否则读取中断会把未扫描的后半桶全部判成删除。
系统验收应写成可测条件。常规更新的可见时间由 Hot 通道验证,历史回改和删除的可见时间由 Warm 通道验证;Worker 在目标端提交后宕机,需要确认重启没有产生重复业务结果;状态库清空后,需要确认 Cold 通道可以重建。边界写清后,技术选型才有可核对的依据。