千万级数据同步与 ETL 实战(三):Hash Wheel、行指纹与删除识别
有些源表只有主键和业务字段,没有更新时间、行版本、删除标记,也拿不到数据库日志。要发现历史记录被修改或删除,完整周期仍需检查全部记录。工程上的问题变成:怎样把一次重扫描拆开,怎样保证同一条记录总落到同一个分区,怎样避免读取中断造成误删除,以及怎样把目标写入量限制在真实变化量附近。
Hash Wheel 负责把全表检查平摊到时间轴,Fingerprint Index 保存上一次看到的状态,Generation 在完整扫描后识别删除。这三个组件必须共同工作。只实现“主键取模加定时查询”,无法形成可恢复的变化检测链路。
一、身份键和内容指纹分开
每条记录先抽象为二元组:
$$ Entry=(IdentityKey, RowFingerprint) $$
IdentityKey 由稳定主键或业务键得到,用于确认记录身份。RowFingerprint 由参与同步的字段计算,用于判断内容变化。两者的判断关系如下:
| 旧状态 | 本轮状态 | 结果 |
|---|---|---|
| 不存在 | 存在 | INSERT |
| 存在 | 指纹不同 | UPDATE |
| 存在 | 指纹相同 | NO_CHANGE |
| 存在 | 完整扫描后未见 | DELETE |
业务键允许修改时,它不再是稳定身份。系统只能把键变更表达成旧键删除和新键新增,除非源端另有不可变 ID 或变更日志能够关联前后值。这个限制需要进入数据集契约。
身份键可直接保存原主键,也可以保存带命名空间的摘要:
$$ IdentityKey=H(datasetId \parallel canonical(primaryKey)) $$
跨数据集使用同一主键时,datasetId 防止状态键冲突。原主键是否保留取决于审计和回查需求;只有摘要会增加碰撞处理与问题定位成本。
二、Hash 前先定义 Canonicalization
指纹错误通常源于编码规则,而非 Hash 函数。以下值在某些业务语义下相同,直接转字符串却会得到不同字节:
Decimal: 1、1.0、1.00
Time: 2026-07-19 10:00:00、2026-07-19T02:00:00Z
JSON: {"a":1,"b":2}、{"b":2,"a":1}
String: Unicode 组合字符与预组合字符每次编码口径调整都要提升 canonicalizationVersion,旧指纹只按旧口径解释。字段规则包括:
| 类型 | 编码规则 |
|---|---|
| NULL | 独立类型标记,和空字符串分开 |
| String | UTF-8,Unicode NFC;Trim 由字段配置决定 |
| Decimal | 固定 Scale,使用十进制定点文本或 unscaled value |
| Timestamp | 转为统一时区和固定精度 |
| Boolean | 单字节 0 或 1 |
| JSON Object | 递归排序 Key,数字规则固定 |
| Array | 保留顺序;若业务视为集合,需单独声明排序规则 |
| Missing | 使用缺失标记,和显式 NULL 分开 |
字段不能只靠分隔符拼接。("a", "bc") 和 ("ab", "c") 在不严谨的拼接规则中可能产生相同字节。长度前缀可以消除边界歧义:
VALUE | fieldType | byteLength | bytes
NULL | fieldTypeJava 编码器可以使用确定的二进制格式:
public final class CanonicalEncoder {
private static final byte NULL = 0;
private static final byte VALUE = 1;
public static void writeString(
DataOutput out,
String value
) throws IOException {
if (value == null) {
out.writeByte(NULL);
return;
}
out.writeByte(VALUE);
String normalized = Normalizer.normalize(
value,
Normalizer.Form.NFC
);
byte[] bytes = normalized.getBytes(StandardCharsets.UTF_8);
out.writeInt(bytes.length);
out.write(bytes);
}
public static void writeDecimal(
DataOutput out,
BigDecimal value,
int scale
) throws IOException {
if (value == null) {
out.writeByte(NULL);
return;
}
out.writeByte(VALUE);
BigDecimal fixed = value.setScale(
scale,
RoundingMode.UNNECESSARY
);
writeString(out, fixed.toPlainString());
}
}RoundingMode.UNNECESSARY 会在精度不符合契约时失败,避免同步过程悄悄改变数值。若业务允许舍入,应在转换规则里明确舍入模式,并让原值和转换错误进入审计记录。
规范化版本发生变化时,旧指纹与新指纹不可直接比较。任务需要进入新 Epoch,重建状态或同时保存旧、新两套指纹直到迁移完成。
三、128 位指纹的碰撞概率
行级快速比较可以使用 xxHash3 128-bit 或 MurmurHash3 128-bit;跨系统强校验可使用 BLAKE3 或 SHA-256。非加密 Hash 适合本地变化检测,不适合作为防篡改凭证。
对均匀的 $b$ 位 Hash,保存 $n$ 个不同输入时,至少一次碰撞的生日界近似为:
$$ p \approx 1-\exp\left(-\frac{n(n-1)}{2^{b+1}}\right) $$
当 $n^2 \ll 2^b$ 时可化简为:
$$ p \approx \frac{n^2}{2^{b+1}} $$
取 $n=10^7$、$b=128$:
$$ p \approx \frac{10^{14}}{2^{129}} \approx 1.47\times10^{-25} $$
该概率足以满足常规数据同步的行级去重。安全审计、跨组织对账或高损失业务可以同时保存快速指纹和 SHA-256,并在桶摘要不一致时使用强摘要复核。概率很低不等于数学上为零,系统文档应保留碰撞这一失败边界。
低基数字段的无密钥摘要可能被枚举。例如状态字段只有几个候选值,攻击者可以预计算结果。摘要需要跨信任边界传输时,可使用 HMAC-SHA-256(secret, datasetId || canonicalRow);密钥轮换会改变全部摘要,需要独立版本号和双版本迁移。
四、稳定分桶把扫描平摊到周期
定义桶号:
$$ bucketId=hash(IdentityKey)\bmod M $$
M 在任务版本中固定。1000 万行、1024 个桶时,均匀分布下每桶期望行数为:
$$ E[rowsPerBucket]=\frac{10,000,000}{1024}\approx9766 $$
每分钟扫描 16 个桶,完整轮转约需:
$$ \frac{1024}{16}=64\text{ 分钟} $$
平均值无法描述热点。租户型主键、带前缀编号和历史空洞都可能造成倾斜。上线前要统计每桶行数、字节数和扫描耗时的分位数。调度成本应基于实际耗时和字节量,不能只数桶。
源表可改造时,可以增加生成列并建立 (sync_bucket, primary_key) 索引。CRC32 可用于分桶,但不用于内容强校验。分桶表达式、字符编码和符号位处理必须跨数据库与 Agent 一致。
ALTER TABLE source_table
ADD COLUMN sync_bucket INT
GENERATED ALWAYS AS (
MOD(CRC32(CAST(id AS CHAR)), 1024)
) STORED;
CREATE INDEX idx_sync_bucket_id
ON source_table(sync_bucket, id);源表不可改造时有三种退路。数值主键可按范围分片,分布倾斜时根据预采样分位点生成近似等量区间。随机键可在应用层 Hash,但源端仍要读取候选范围,无法减少数据库扫描量。若查询条件写成 MOD(HASH(id), M)=? 且没有函数索引,数据库可能每个桶都全表计算,1024 个桶会把一次全扫放大为大量重复工作。
M 变更会让多数记录重新分桶。迁移时应建立新的 Wheel Epoch,双写或重建新状态,校验后切换,不能直接修改配置后沿用旧 Generation。
五、Fingerprint Index 的状态布局
索引保存上一轮看到的记录:
Key = datasetId | wheelEpoch | bucketId | identityKey
Value = fingerprint | lastSeenGeneration | sourceVersion | schemaId千万级本地状态适合 RocksDB、LMDB、SQLite 等嵌入式存储,也可以放中央 PostgreSQL。选择依据包括单点恢复方式、并发写模型、磁盘写放大、备份时间和是否需要跨 Worker 共享。把全部状态放 Redis 前应计算内存编码开销、持久化峰值和淘汰风险。
以 1000 万条状态估算,裸字段约 56 至 72 字节每条,总计 560MB 至 720MB。KV 的 Key 前缀、索引、WAL、Block、Bloom Filter 与 Compaction 会把实际磁盘占用推高到约 1.5GB 至 4GB。长主键直接进入 Key 时还会增加索引和缓存压力,可以保存 128 位身份摘要,并把原键放入变化 Batch 或单独的冲突记录。
状态写入与变化 Batch 要有恢复协议。若先覆盖旧指纹,再在 Batch 落盘前宕机,恢复扫描会认为记录没有变化,导致更新丢失。可采用 Write Batch 同时写“候选状态”和 Batch Manifest,Sink 提交后再提升为 committed generation;也可以让状态更新由不可变 Batch 重放生成。第四篇会展开提交顺序。
六、Generation 只能在完整桶扫描后提交
扫描桶 $b$ 时,为本次运行分配:
$$ g_{new}=g_{committed}+1 $$
每读到一条记录,计算身份和指纹,并把 lastSeenGeneration 写为 $g_{new}$。本地不存在则产生 INSERT,指纹变化则产生 UPDATE,指纹一致只刷新本轮可见标记。
桶完整读取、差异 Batch 持久化、目标端写入成功后,仍满足以下条件的旧状态才可产生 DELETE:
$$ lastSeenGeneration < g_{new} $$
读取到一半时禁止执行这一步。否则后半桶尚未读取的记录会全部变成墓碑。
分页游标、候选 Generation 和扫描状态必须一起恢复。失败后可以从桶起点重扫,代价较高但语义简单;也可以从已持久化页游标继续,但此前页面的 lastSeenGeneration 和 Batch 必须已经可靠落盘。任何“只保存 lastId、不保存运行 Generation”的实现都难以区分本轮数据和上轮残留。
删除事件使用 Tombstone,包含身份键、源版本、删除发现时间和扫描证据。目标端是否物理删除由保留策略决定。对于软删除、审计和下游迟到事件,保留墓碑一段时间通常更容易恢复。
七、热点桶采用 EWMA 调度
固定轮询能给每个桶同样的次数,却不能给出同样的变化发现延迟。活跃桶需要更频繁,历史冷桶可以降低频率,但每个桶仍有强制陈旧上限。
设桶本轮变化率为 $x_t$,指数移动平均为:
$$ \lambda_t=\alpha x_t+(1-\alpha)\lambda_{t-1} $$
α 越大,对近期突增响应越快,也更容易受短期噪声影响。调度优先级可以组合陈旧度、SLO、变化率和扫描成本:
$$ priority_i= \frac{(age_i/SLO_i)^\gamma(\lambda_i+\varepsilon)}{cost_i} $$
其中 $\varepsilon>0$ 防止长期无变化的桶优先级归零,$\gamma>1$ 让接近超期的桶快速提升。再加硬约束:
$$ age_i \ge maxStaleness_i \Rightarrow \text{必须调度} $$
两个候选桶的输入参数如下:
| 桶 | 已等待 | SLO | 变化率 | 扫描成本 |
|---|---|---|---|---|
| A | 10 分钟 | 20 分钟 | 2% | 1 秒 |
| B | 50 分钟 | 60 分钟 | 0.1% | 0.2 秒 |
公式参数不同会改变排序,调度器需要输出每次选择的分项得分,便于解释冷桶为何获得执行机会。
EWMA 状态也要版本化。调度算法升级后先离线回放历史指标,检查完整轮转时间、热点延迟和源端 QPS,再灰度发布。只优化平均延迟可能造成少数大桶长期饥饿。
八、扫描速度由源端和目标端共同约束
假设目标完整周期为 $C$,数据量为 $N$,平均每行投影字节为 $s$,则源端持续带宽下限约为:
$$ B_{source}\ge\frac{N\times s}{C} $$
1000 万行、每行投影 200 字节、1 小时一轮,持续读取约 0.56MB/s。网络数字不大,但数据库还要付出索引遍历、行可见性判断和随机页读取成本。若同步投影包含大字段,带宽和 Buffer 压力会显著增加。变化检测阶段可只读身份键、参与指纹的列和版本字段,确认变化后再回表读取大对象,但两次读取之间可能发生并发修改,需要版本条件或一致性策略。
Chunk 大小可以按目标耗时反馈:
$$ nextSize=clip\left( currentSize\times\frac{targetLatency}{actualLatency}, minSize, maxSize \right) $$
实际控制器还要限制单次变化幅度,并同时约束行数、字节数和查询超时。一条超大记录不能让整个批次无限膨胀,应进入大记录通道或明确失败。
九、测试关注不变量
示例数据要覆盖 NULL 与空字符串、不同 Decimal Scale、Unicode 组合字符、JSON Key 顺序、复合键、Hash 桶倾斜和同键并发修改。属性测试适合验证:相同规范化值产生相同指纹;摘要合并满足结合律;未完整扫描的桶不会产生 DELETE;提交的 Generation 单调增加;重复执行相同 Batch 不改变结果。
故障注入应落在状态边界:扫描第 3 页断网、候选指纹写入后宕机、Sink 提交后 Checkpoint 失败、租约过期后旧 Worker 继续写、Schema 在桶中途变化、分桶索引失效导致查询退化。验收结果既要比较目标数据,也要检查源端扫描行数、完整轮转时间和状态库磁盘水位。
Hash Wheel 的价值来自稳定分区和可恢复状态。它没有让 $O(N)$ 扫描消失,却把一次不可控的大任务拆成可限速、可暂停、可重试的小任务,并把网络与目标写入压缩到变化集合附近。第五篇会在此状态之上增加桶摘要、Merkle Tree 和 IBLT,减少跨节点对比和差异定位成本。