千万级数据同步与 ETL 实战(四):Micro Batch、幂等写入与故障恢复
增量检测得到一组 INSERT、UPDATE 和 DELETE 后,数据仍未安全到达目标端。读取进程可能宕机,网络可能在响应返回前断开,目标库可能提交成功却没有被调用方确认,两个 Worker 也可能短时间内同时持有同一分片。系统能否恢复,取决于批次是否可重放、写入是否幂等、Checkpoint 何时推进,以及旧租约持有者能否继续提交。
Micro Batch 的作用不只是凑够若干行再发送。它是抽取结果的持久化边界,也是源端读取和目标端提交之间的恢复凭证。
一、Batch 创建后保持不可变
一个批次由数据文件和 Manifest 组成。数据文件可使用 Arrow、Parquet 或有明确 Schema 的二进制格式;Manifest 描述这批数据从哪里来、覆盖哪个游标区间、包含多少变化,以及如何验证内容。
{
"batchId": "01J2SYNTHETIC8TQ5F6X8R3",
"taskId": "customer-sync",
"taskVersion": 12,
"partitionId": "bucket-0037",
"epoch": 42,
"schemaId": "customer-v7",
"cursorFrom": "2026-07-19T10:00:00.000Z|80000",
"cursorTo": "2026-07-19T10:01:00.000Z|84321",
"recordCount": 5000,
"insertCount": 120,
"updateCount": 47,
"deleteCount": 3,
"dataUri": "batch://customer-sync/42/bucket-0037/part-0001.parquet",
"checksum": "sha256:...",
"createdAt": "2026-07-19T10:01:12Z"
}batchId 一经生成就对应固定字节。重试只能重复发送该批次,不能复用同一个 ID 重新查询源端并覆盖内容。源数据可能在两次查询之间变化,同 ID 不同内容会破坏去重和审计。
批次写入可采用临时对象加原子完成标记:先写数据文件,计算长度和校验和,再写 Manifest,随后将状态改为 READY。本地文件系统可使用同卷原子重命名;对象存储通常没有目录重命名语义,可以用不可变对象 Key 和独立状态记录。看到 READY 前,Sink Worker 不得消费。
Batch 状态表只记录生命周期,不允许修改数据文件。目标端出现永久字段错误时,修复映射后应生成新的任务版本和新 Batch;直接改旧文件会让历史校验和失效。
二、Checkpoint 的提交顺序
安全顺序可以写成六个持久化步骤:
1. 读取源端页面,但不推进已提交游标
2. 规范化并生成 ChangeRecord
3. 持久化不可变 Batch
4. 目标端按幂等键提交 Batch
5. 记录目标提交凭证
6. CAS 推进 Checkpoint该顺序允许步骤 4 成功、步骤 6 失败。恢复后会重复写目标端,因此第 4 步必须幂等。反过来,若先推进 Checkpoint 再写目标端,写入失败后源端游标已越过该范围,系统没有可靠材料可以补写。
Checkpoint 至少包含 taskId、partitionId、epoch、游标类型、游标值、lastBatchId 和乐观锁版本。游标编码也要带版本。二元 Watermark、HTTP Cursor、Hash Wheel Generation 和文件版本不能用同一套无类型字符串随意解释。
三、幂等键要覆盖业务身份和版本
只按 batchId 去重能够阻止同一批次重复执行,却不能处理重叠窗口、Hot/Warm 双通道或重建 Epoch 发现的同一变化。目标端还需要记录级幂等条件。
常见设计为:
batch dedupe key = (taskId, batchId)
record identity = (datasetId, identityKey)
record version = sourceVersion目标端先登记批次。如果 batchId 已处于 COMMITTED,直接返回原提交结果;如果批次处于 WRITING 且租约未过期,避免并发执行;租约过期后可以接管。记录 Upsert 只接受不旧于当前版本的数据。
INSERT INTO target_record(
dataset_id, identity_key, payload, source_version, batch_id
)
VALUES (
:dataset_id, :identity_key, :payload, :source_version, :batch_id
)
ON CONFLICT (dataset_id, identity_key) DO UPDATE
SET payload = EXCLUDED.payload,
source_version = EXCLUDED.source_version,
batch_id = EXCLUDED.batch_id
WHERE target_record.source_version <= EXCLUDED.source_version;这段 SQL 只展示意图。sourceVersion 的编码顺序、数据库 Collation、NULL 处理和并发锁行为都要实测。版本相同但内容不同属于数据契约冲突,应拒绝并报警,不能静默覆盖。
无法得到可比较源版本时,可以保证“相同 Batch 重放结果不变”,却无法在多个通道并发时判断旧值和新值。可选办法是让同一身份键固定进入同一分区并串行写入,或使用目标端接收顺序作为弱版本,再由定期对账收敛。文档需要把它标成保证降级。
四、Tombstone 让删除可以重放
物理删除不能只执行一次 DELETE FROM target WHERE id=?。重放、下游订阅、缓存和审计都需要知道某个身份键在某个版本被删除。变化记录应携带 Tombstone:
{
"identityKey": "customer:10086",
"operation": "DELETE",
"sourceVersion": "wheel:42:37:18",
"detectedAt": "2026-07-19T11:08:00Z",
"evidence": {
"bucketId": 37,
"generation": 18,
"scanComplete": true
}
}Tombstone 版本必须高于目标端当前值,且只能由完整日志事件、明确软删字段或已提交的完整桶扫描产生。扫描失败时没有足够证据证明记录消失。
墓碑保留时间取决于下游迟到上限和重建周期。过早清理会让迟到旧值重新创建记录。保留期内若收到更旧的 UPSERT,应拒绝;收到有证据的新版本,则允许恢复记录。
五、分片租约还需要 Fencing Token
租约解决 Worker 宕机后的接管,但单靠 leaseUntil 无法阻止旧 Worker。旧 Worker 可能经历长时间 GC 或网络隔离,租约已经过期,新 Worker 完成接管;旧 Worker 恢复后仍认为自己有权提交。
租约每次成功领取都增加单调 fencingToken:
UPDATE etl_partition_lease
SET owner_id = :owner,
lease_until = :lease_until,
fencing_token = fencing_token + 1,
version = version + 1
WHERE task_id = :task_id
AND partition_id = :partition_id
AND (
owner_id IS NULL
OR lease_until < CURRENT_TIMESTAMP
OR owner_id = :owner
);Batch、Checkpoint 和目标端批次登记都携带该 Token。持久化层拒绝小于当前 Token 的写入。租约只说明“当前时间是否仍有效”,Fencing Token 负责隔离恢复后的旧持有者。
逻辑分片数可以大于 Worker 数,便于负载均衡。并发上限由源端连接数、查询延迟、磁盘写入和 Sink 吞吐共同决定。1000 万行任务并不需要大量线程;无节制并发会把业务数据库和本地 Compaction 同时推到高水位。
六、动态 Chunk 受行数、字节和延迟共同限制
固定 5000 行无法覆盖窄表和大字段表。控制器可根据目标延迟调整下一批行数:
$$ n_{t+1}=clip\left( n_t\times\frac{L_{target}}{L_{actual}}, n_{min},n_{max} \right) $$
为避免一次测量造成剧烈波动,再限制变化比例,例如每轮只允许扩大 25% 或缩小 50%。同时设置 maxChunkBytes、查询超时和单记录上限。以下任一条件成立时结束当前 Chunk:
rows >= maxChunkRows
or bytes >= maxChunkBytes
or readTime >= maxChunkLatency一条记录大于 maxChunkBytes 时不能无限重试。系统可以让它独占一个 Batch,或转入大对象通道;超过 Sink 硬限制则进入 Dead Letter,并记录字段长度和规则版本,不能在日志中打印敏感正文。
七、背压沿链路向源端传播
Sink 变慢时继续高速抽取,会把压力转移到本地磁盘和内存。需要同时观察:
pending_batch_count
pending_batch_bytes
oldest_pending_batch_age
sink_write_latency
sink_error_rate
local_disk_usage
source_read_latency
source_connection_usage控制策略可以分层。队列字节超过软水位时降低新读取并发;磁盘达到硬水位时暂停低优先级 Warm 和 Cold 任务;Sink 错误率连续升高时打开熔断,只保留探测请求;源端查询延迟上升时缩小 Chunk 和令牌桶速率。
反压反馈要有滞回区间。达到 80% 降速、刚降到 79% 就恢复,会导致频繁振荡。可在 80% 触发暂停,降到 60% 后分阶段恢复。恢复速度要慢于降速速度,避免下游刚恢复就被积压流量再次压垮。
八、Schema 变更不能混入执行中的 Batch
Worker 启动时取得一份带 schemaId 的发布快照,当前 Batch 完成前不刷新。可空字段新增通过后续任务版本引入;删除参与同步的列或改变身份键时,任务先停在 Batch 边界,再重建受影响的目标结构和指纹状态。
Schema 状态可以使用:
DETECTED
→ PAUSED
→ MIGRATION_PLANNED
→ VALIDATED
→ RESUMED同一 Batch 内不能混用两个 Schema。Sink 先验证 Manifest 中的 schemaId,再读取数据。任务回滚时要确认目标端仍能解释旧 Schema,也要保留对应转换代码和字典。
兼容性判定不能只看字段名。Decimal Scale、时区、字符排序规则、JSON 数字类型和缺失值语义都可能改变指纹或目标值。Schema Registry 应保存逻辑类型和规范化策略,物理数据库类型只是其中一部分。
九、逐个故障点定义恢复动作
故障恢复要从提交顺序推导:
| 故障位置 | 可见状态 | 恢复动作 |
|---|---|---|
| Source 读取失败 | 无 READY Batch | 不推进游标,从页游标重试 |
| Batch 落盘失败 | 临时文件或 CREATING | 清理未完成对象,重新读取 |
| Sink 写入失败 | READY Batch 保留 | 不回源,重放同一 Batch |
| Sink 成功、响应丢失 | 结果未知 | 查询批次登记或幂等重写 |
| Checkpoint CAS 失败 | Sink 可能已提交 | 停止旧 Worker,加载新状态 |
| Worker 租约过期 | 新 Worker 可能接管 | 依赖 Fencing Token 拒绝旧写 |
| 状态库损坏 | 增量比较不可用 | 暂停相关分片,从 Cold 基线重建 |
指数退避要带随机抖动,避免大量分片同时恢复造成重试峰值。永久错误与暂时错误分开:网络超时、限流和主库切换通常可重试;字段无法转换、权限撤销和 Schema 不兼容需要暂停并人工处理。
Dead Letter 不是无限堆放区。每条记录要有错误分类、任务版本、Schema、批次、首次和末次失败时间、重放次数及脱敏后的诊断信息。修复后通过新版本重放,并验证 Checkpoint 没有越过无法接受的数据范围。
十、正确性测试以状态转换为中心
端到端测试准备合成数据,并在每个持久化步骤后强制退出进程。重新启动后,验证目标端收敛状态、批次状态、Checkpoint、Generation 和租约 Token。关键不变量包括:已提交 Checkpoint 对应的所有 Batch 均已在目标端提交;未 READY 的 Batch 不可消费;同一 Batch 重放不改变业务结果;小 Token 无法覆盖大 Token;未完成桶扫描不产生墓碑。
吞吐压测全程保留状态不变量检查。高吞吐结果若来自提前推进 Checkpoint、丢弃错误记录或关闭版本比较,没有参考价值。生产演练应覆盖磁盘高水位、目标端持续超时、连接池耗尽和状态库恢复,确认系统会降速和停在可解释的位置。
“Exactly Once”不能作为一句验收口号。微批链路的可验证保证可以写成:源端可能重复读取;READY Batch 可重复投递;目标端按批次和记录版本幂等;Checkpoint 只在目标提交后推进。每一条都能通过故障注入检查,系统的正确性由这些不变量共同构成。