采集任务断点续传实现方案:先为任务选定稳定游标并保存已提交检查点,再按小批次读取数据;每批在同一个数据库事务中完成结果幂等写入和检查点推进。执行器中断后,由新执行器从最后一次已提交的检查点继续,并通过租约、所有权校验和重试机制处理并发与故障。

采集任务断点续传实现方案:检查点、幂等写入与故障恢复实操

采集任务断点续传是指采集进程中断、机器重启或任务被接管后,依据持久化的已完成进度继续采集,而不必从头开始。这里的“断点”不是当前正在读取的位置,而是结果已经成功提交的位置;未提交的批次允许重新执行。

一、第一步:确定检查点和采集上界

1. 选择稳定的进度单位

优先使用单调递增且唯一的源数据 ID、上游提供的稳定分页游标,或由更新时间与唯一 ID 组成的复合游标。普通页码会随新增、删除和排序变化而漂移,除非上游提供固定快照,否则不宜作为检查点。

使用 (updated_at, id) 时,查询条件和排序必须一致:

WHERE updated_at > :last_time
   OR (updated_at = :last_time AND id > :last_id)
ORDER BY updated_at ASC, id ASC
LIMIT :batch_size

该游标适用于源端能够稳定返回排序结果的场景。若旧记录可能在采集过程中更新并移动到游标之后,应结合固定窗口、下一轮重叠扫描或源端变更日志处理,不能只依赖一次遍历。

2. 固定本轮范围

任务启动时记录可执行的采集上界,例如 snapshot_max_id、时间窗口结束点或上游快照令牌。本轮只处理该范围内的数据,新数据留给下一轮。时间上界并不等于数据库快照;若源端存在延迟写入,下一轮仍需设置回溯窗口并依靠幂等写入去重。

二、第二步:建立任务表与结果表

1. 持久化任务状态和检查点

下面以支持行锁和 JSON 字段的关系型数据库为例;数据类型及领取语法需按实际数据库调整。

CREATE TABLE collection_task (
  id BIGINT PRIMARY KEY,
  source_key VARCHAR(128) NOT NULL,
  status VARCHAR(20) NOT NULL,
  checkpoint_json JSON NOT NULL,
  snapshot_json JSON NULL,
  owner_id VARCHAR(64) NULL,
  lease_until TIMESTAMP NULL,
  attempt_count INT NOT NULL DEFAULT 0,
  max_attempts INT NOT NULL DEFAULT 10,
  next_run_at TIMESTAMP NULL,
  last_error TEXT NULL,
  version BIGINT NOT NULL DEFAULT 0,
  created_at TIMESTAMP NOT NULL,
  updated_at TIMESTAMP NOT NULL
);

checkpoint_json保存最后已提交的 ID、复合游标或分页令牌;snapshot_json保存本轮上界。owner_idlease_untilversion用于任务接管与失效执行器隔离。

2. 用唯一约束兜底幂等写入

CREATE TABLE collected_item (
  id BIGINT PRIMARY KEY,
  source_key VARCHAR(128) NOT NULL,
  source_item_id VARCHAR(128) NOT NULL,
  payload JSON NOT NULL,
  collected_at TIMESTAMP NOT NULL,
  UNIQUE (source_key, source_item_id)
);

对于会更新的记录,按业务规则使用 UPSERT;对于不可变记录,可在唯一键冲突时忽略重复写入。不能仅靠应用程序先查询、再插入实现去重。

三、第三步:原子领取任务并维护租约

1. 在短事务中领取任务

多个执行器并发运行时,可在事务内使用 SELECT ... FOR UPDATE SKIP LOCKED锁定一条待处理任务,再写入执行器标识、租约截止时间并递增版本号。待处理或重试任务应满足 next_run_at已到期;运行中任务仅在租约过期后允许接管。领取事务提交后再请求上游接口,不要在持有行锁时执行网络采集。

2. 续租并阻止旧执行器提交

心跳可按租约时长的约三分之一执行,例如租约 90 秒、每 30 秒续租。续租时校验任务状态、owner_id、当前版本及未过期的租约。批次提交时也必须在同一事务内校验这些条件;条件更新影响行数不是 1 时,回滚结果写入并停止该执行器。租约判断尽量使用数据库时间,减少机器时钟漂移的影响。

四、第四步:事务性提交结果与检查点

1. 按批次执行

  1. 读取任务的检查点、上界及当前所有权版本。
  2. 从检查点之后拉取一批数据,完成解析和校验。
  3. 开启数据库事务,按业务唯一键写入结果。
  4. 在事务内通过所有权、租约和版本条件更新检查点;影响行数不是 1 则回滚。
  5. 提交事务,再处理下一批;确认范围已采完后,按同样的所有权条件将任务标记为成功。

例如每批先处理 100 条,再根据源接口延迟、事务耗时和恢复成本调整批次大小。检查点必须取该批最后一条已纳入本次提交范围的记录游标,不能在读取完成时提前推进。

2. 事务伪代码

task = load_task(task_id)
batch = source.fetch_after(task.checkpoint, task.snapshot, limit=100)

if batch.is_empty():
    mark_success_if_owned(task.id, worker_id, task.version)
else:
    begin_transaction()
    try:
        upsert_items(batch.items)
        affected = update_checkpoint_if_owned_and_lease_valid(
            task.id, worker_id, task.version,
            task.checkpoint, cursor_of(batch.last_item)
        )
        if affected != 1:
            raise LeaseLostError()
        commit()
    except:
        rollback()
        raise

这里假设结果表和任务表位于同一个支持事务的数据库。事务回滚时结果和检查点都不生效;事务提交后即使执行器未收到成功响应,重放批次也会由唯一约束和 UPSERT 处理。

五、第五步:配置重试与恢复流程

1. 区分错误类型

网络超时、限流和临时数据库故障可按指数退避重试,例如 min(5秒 × 2^attempt, 30分钟) + jitter。鉴权失效、参数错误或持续无法解析的数据不应无限重试,应记录脱敏错误并暂停或转入人工处理。达到最大次数后将任务标记为失败;单条异常数据可按业务规则隔离,但不能未经记录就跳过并推进检查点。

五、第五步:配置重试与恢复流程

2. 从已提交位置恢复

扫描租约过期的运行中任务,将其转为可重试状态,或允许领取逻辑直接接管。新执行器读取数据库中的最后已提交检查点,从其后继续读取。未提交批次可能再次执行,这是“至少一次执行 + 幂等写入”的预期行为。

六、第六步:处理数据变化和外部存储

1. 增量更新、删除与文件变化

对于持续更新的源数据,可使用 (updated_at, id)和固定时间窗口,并在下一轮适度回溯,例如从上轮结束时间前 5 分钟重新读取;回溯长度应依据源端最大可接受延迟设定。物理删除需依靠删除事件、逻辑删除标记或定期对账发现。文件采集除偏移量外,还应保存文件标识、大小和摘要;恢复前确认文件未被替换。普通压缩文件可能无法按字节偏移随机续读,可改用分片或独立压缩块。

2. 外部结果不能加入本地事务时

写入搜索引擎、对象存储或第三方接口时,本地数据库事务不能保证外部写入与检查点同时成功。可先在同一事务中保存待投递内容、outbox 记录和源端检查点,再由投递器按稳定幂等键重试外部写入。此时源端检查点仅表示内容已可靠暂存,不表示外部系统已完成写入;任务完成条件应另行检查 outbox 是否全部投递成功。外部接口不支持幂等或结果查询时,还需对账与补偿。

七、第七步:开展故障注入与一致性验证

1. 验证基本恢复路径

  1. 准备 1000 条具有稳定唯一 ID 的测试数据,设置每批 100 条。
  2. 在第 3 批提交前强制终止执行器,确认检查点仍停留在第 2 批已提交位置。
  3. 重新启动或由另一执行器接管任务,等待任务完成。
  4. 核对目标端记录数和唯一 ID 数均为 1000,并检查内容摘要及最终检查点。

2. 覆盖事务与租约边界

分别在拉取后、事务提交前、事务提交成功但执行器收到响应前终止进程;再模拟旧执行器暂停至租约过期后恢复、两个执行器竞争同一任务、上游超时和重复返回数据。预期结果是未提交批次可重放、已提交结果不重复、失去租约的执行器无法推进检查点。

3. 核对上线指标

检查源端与目标端的分区数量、业务唯一键、关键字段摘要和最终检查点;使用 outbox 时,还要检查待投递数量。上线后持续监控检查点推进速度、任务延迟、租约接管次数、批次耗时、重试率和隔离数据量。

八、上线注意事项与总结

将检查点保存在具备事务和备份能力的持久化存储中;Redis 可辅助调度,但不宜作为唯一进度依据。取消任务时在批次边界停止并保留已提交检查点,不要误标为成功。日志应避免记录密钥、Cookie 和完整敏感响应。

总结:这套采集任务断点续传实现方案的最小闭环是“稳定游标 + 持久化检查点 + 结果唯一约束 + 同事务提交 + 租约接管”。上线前用故障注入验证每个提交边界,并用数量、唯一键和内容对账确认既能恢复,也不会漏采。

常见问题

断点续传应该保存页码还是最后一条数据的 ID?

优先保存最后一条已提交记录的稳定 ID 或复合游标。普通页码可能因源数据新增、删除或排序变化而漂移;只有固定快照下的页码才适合作为检查点。

为什么结果写入和检查点更新要放在同一个事务中?

同一数据库事务可避免检查点已推进、结果却未落库导致的漏采,也能在提交失败时让整批安全重放。即使发生结果已提交但响应丢失,唯一约束和 UPSERT 仍需负责处理重复执行。

租约过期后,如何防止旧执行器覆盖新进度?

接管时更新 owner_id 和版本号;每次续租、提交检查点及完成任务时,校验当前所有者、版本和租约。条件更新失败时回滚本批事务并停止旧执行器。

使用 Redis 保存采集进度是否足够?

不建议将 Redis 临时键作为唯一进度依据。关键检查点应存放在支持持久化、事务和备份的存储中,并与结果提交保持一致。

结果写入外部系统时,什么时候算任务完成?

可先将待投递内容、outbox 记录和源端检查点提交到本地数据库,再重试投递。只有外部投递成功且待投递记录清零后,才能按外部交付目标认定任务完成。

如何验证断点续传可以上线?

至少覆盖提交前后进程终止、租约过期接管、双执行器竞争、上游超时和重复数据,并核对最终检查点、数量、唯一键、内容摘要及待投递记录。

总结

采集任务断点续传是依据持久化的已提交检查点,在中断后继续执行。落地时使用稳定游标和固定采集范围,在同一数据库事务中完成结果幂等写入与检查点推进,再用租约、版本校验和有限重试处理执行器故障。外部存储需配合 outbox、幂等投递与对账;上线前通过故障注入验证恢复和数据一致性。