← 返回

千万级数据同步与 ETL 实战(三):Hash Wheel、行指纹与删除识别

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

有些源表只有主键和业务字段,没有更新时间、行版本、删除标记,也拿不到数据库日志。要发现历史记录被修改或删除,完整周期仍需检查全部记录。工程上的问题变成:怎样把一次重扫描拆开,怎样保证同一条记录总落到同一个分区,怎样避免读取中断造成误删除,以及怎样把目标写入量限制在真实变化量附近。

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独立类型标记,和空字符串分开
StringUTF-8,Unicode NFC;Trim 由字段配置决定
Decimal固定 Scale,使用十进制定点文本或 unscaled value
Timestamp转为统一时区和固定精度
Boolean单字节 01
JSON Object递归排序 Key,数字规则固定
Array保留顺序;若业务视为集合,需单独声明排序规则
Missing使用缺失标记,和显式 NULL 分开

字段不能只靠分隔符拼接。("a", "bc")("ab", "c") 在不严谨的拼接规则中可能产生相同字节。长度前缀可以消除边界歧义:

VALUE | fieldType | byteLength | bytes
NULL  | fieldType

Java 编码器可以使用确定的二进制格式:

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{ 分钟} $$

flowchart LR K["IdentityKey"] --> H["固定 Hash"] H --> M["mod 1024"] M --> B["稳定 Bucket"] B --> P["Keyset 分页扫描"] P --> F["Fingerprint Index 对比"] F --> D["变化 Batch / Generation 提交"]

平均值无法描述热点。租户型主键、带前缀编号和历史空洞都可能造成倾斜。上线前要统计每桶行数、字节数和扫描耗时的分位数。调度成本应基于实际耗时和字节量,不能只数桶。

源表可改造时,可以增加生成列并建立 (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} $$

读取到一半时禁止执行这一步。否则后半桶尚未读取的记录会全部变成墓碑。

stateDiagram-v2 [*] --> CREATED CREATED --> SCANNING SCANNING --> SCANNED: 完整读到页尾 SCANNING --> FAILED: 读取或落盘失败 SCANNED --> DIFF_COMPUTED DIFF_COMPUTED --> WRITING WRITING --> COMMITTED: Sink 与 Generation 提交 WRITING --> FAILED: 写入失败 FAILED --> SCANNING: 从已保存游标重试

分页游标、候选 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变化率扫描成本
A10 分钟20 分钟2%1 秒
B50 分钟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,减少跨节点对比和差异定位成本。

参考资料