← 返回

千万级数据同步与 ETL 实战(一):边界、降级策略与整体架构

系列导航:1 | 2 | 3 | 4 | 5 | 6

数据同步项目开工时,团队常先讨论 Kafka、Flink 或 Debezium。这个顺序容易忽略一个基础问题:源系统究竟能提供什么。很多企业数据源只有只读账号、分页接口或定期投递的文件。数据库日志权限拿不到,表里也未必有可靠的更新时间。此时直接套用 CDC 架构,实施阶段会在权限、版本、网络和运维能力上不断返工。

本文先给系统划定边界。目标数据集以单表或单集合不超过 1000 万行为基准,变更速率不高,业务接受几十秒到数小时的延迟。系统需要发现新增、修改和删除,支持断点恢复,并限制对源库的影响。每秒数十万事件、复杂事件时间计算、强实时交易链路不在这套微批架构的处理范围内,第六篇会单独讨论真正的流式系统。

一、先盘点数据源能力

同一个“查询数据”接口,能提供的增量能力差别很大。建任务前要探测并固化以下信息:

能力需要确认的内容缺失后的影响
变更日志Binlog、WAL、LSN、SCN、日志保留时间无法按日志位置持续消费
版本字段单调版本号、精度、是否覆盖删除只能依赖时间或扫描
更新时间是否有索引、是否随每次修改更新历史回改可能漏读
稳定键单主键、复合业务键、键是否会变化无法可靠识别同一条记录
一致性快照隔离级别、快照时长、日志衔接点首次全量与增量可能断层
分页方式游标、Keyset、页码、排序稳定性并发写入时可能漏读或重复
删除信息删除日志、软删字段、墓碑记录需要周期扫描才能识别删除
源端改造能否加索引、生成列、触发器分桶扫描可能退化成全表计算

能力探测结果进入数据集版本,不能只留在部署人员的经验里。任务发布后如果主键、时区或更新时间精度发生变化,旧 Checkpoint 可能失效。调度器需要暂停任务,完成迁移验证后再恢复。

建议把增量策略设计成一条降级链:

flowchart LR A["读取源端能力"] --> B{"可消费变更日志"} B -->|是| C["CDC / 日志增量"] B -->|否| D{"有版本号或更新时间"} D -->|是| E["Version / Watermark"] D -->|否| F{"数据只追加"} F -->|是| G["自增键增量"] F -->|否| H["Hash Wheel 分桶扫描"] H --> I["全量指纹校验兜底"]

降级需要记录原因。例如数据库支持 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 重启后既不知道应该执行哪个版本,也不知道哪个批次已经对外生效。

flowchart TB subgraph CP["控制面"] CR["连接器与数据集注册"] SR["Schema / 任务版本"] SP["调度与分片规划"] CL["Checkpoint / 租约 / 重试"] CR --> SR --> SP --> CL end subgraph DP["数据面"] RD["Source Reader"] --> PS["Partition Splitter"] PS --> ID["Increment Detector"] ID --> NM["Canonicalizer / Transformer"] NM --> MB["Immutable Micro Batch"] MB --> SW["Sink Writer"] end SP --> PS CL <--> ID SW --> CL

控制面发布的是不可变任务版本,至少包括数据源能力快照、字段投影、身份键、规范化规则、增量模式、分片算法、目标写入语义和 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
}

identityKeyfingerprint 不能合并。前者回答“是否为同一条记录”,后者回答“内容是否改变”。如果把整行摘要当作身份,修改一列会被解释成删除旧记录再新增记录,外键关系、审计链和下游聚合都会受到影响。

sourceVersion 也不能默认使用抽取时间。抽取时间只表示 Worker 何时读到记录,无法阻止迟到的旧值覆盖新值。这个字段必须来自源端可验证的先后关系;第二篇会结合 Watermark 讨论具体编码。只有行指纹可用时,目标端具备幂等性,却未必能比较新旧顺序,需要由同一分区串行处理,或在冲突时回源核对。

五、三条校验通道

只有一条 Watermark 通道时,人工回写旧时间、软删字段遗漏更新、应用时钟回拨都会形成长期差异。可以按不同成本设置三条通道:

flowchart LR S["同一数据集"] --> H["Hot:日志或 Watermark"] S --> W["Warm:Hash Wheel"] S --> C["Cold:全量摘要"] H --> T["秒到分钟发现常规变化"] W --> U["分钟到小时发现历史回改和删除"] C --> V["低频重建与全局对账"]

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 通道可以重建。边界写清后,技术选型才有可核对的依据。

参考资料