← 返回

千万级数据同步与 ETL 实战(四):Micro Batch、幂等写入与故障恢复

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

增量检测得到一组 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 不得消费。

stateDiagram-v2 [*] --> CREATING CREATING --> READY: 数据和 Manifest 校验通过 CREATING --> ABORTED: 写入失败 READY --> WRITING: Sink 领取 WRITING --> COMMITTED: 目标端确认 WRITING --> READY: 可重试错误 WRITING --> DEAD: 超过人工处理阈值 COMMITTED --> EXPIRED: 超过重放保留期

Batch 状态表只记录生命周期,不允许修改数据文件。目标端出现永久字段错误时,修复映射后应生成新的任务版本和新 Batch;直接改旧文件会让历史校验和失效。

二、Checkpoint 的提交顺序

安全顺序可以写成六个持久化步骤:

1. 读取源端页面,但不推进已提交游标
2. 规范化并生成 ChangeRecord
3. 持久化不可变 Batch
4. 目标端按幂等键提交 Batch
5. 记录目标提交凭证
6. CAS 推进 Checkpoint

该顺序允许步骤 4 成功、步骤 6 失败。恢复后会重复写目标端,因此第 4 步必须幂等。反过来,若先推进 Checkpoint 再写目标端,写入失败后源端游标已越过该范围,系统没有可靠材料可以补写。

sequenceDiagram participant R as Reader participant B as Batch Store participant S as Sink participant C as Checkpoint Store R->>R: 读取 cursorFrom..cursorTo R->>B: 写入 Batch 和 checksum B-->>R: READY R->>S: write(batchId) S-->>R: commitToken R->>C: CAS(expectedVersion, cursorTo, batchId) alt CAS 成功 C-->>R: committed else CAS 冲突 C-->>R: rejected,停止旧 Worker end

Checkpoint 至少包含 taskIdpartitionIdepoch、游标类型、游标值、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 和令牌桶速率。

flowchart RL S["Sink 延迟 / 错误"] --> Q["Pending Batch 水位"] Q --> W["降低 Worker 并发"] W --> R["减少 Source 读取令牌"] D["磁盘硬水位"] --> P["暂停 Warm / Cold"] P --> W

反压反馈要有滞回区间。达到 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 只在目标提交后推进。每一条都能通过故障注入检查,系统的正确性由这些不变量共同构成。

参考资料