断点续跑与幂等性
---
章节:进阶篇第7章
难度:⭐⭐⭐ 进阶级
阅读时间:20-25 分钟
核心收获:掌握断点续跑和幂等性设计,确保任务失败后可恢复
9.1 理论要点:为什么需要断点续跑
9.1.1 断点续跑的价值
问题:推送失败后,如果重新执行全部步骤,会: - 浪费时间和成本 - 可能产生不同的聚合结果(数据源已更新) - 增加数据源压力
断点续跑的价值:从上次成功的位置继续,不重做已完成步骤。
9.1.2 幂等性的本质
幂等(Idempotency):同一操作执行多次,结果一致,不会产生副作用。
为什么需要幂等: - 手动触发与定时触发同时发生 - 推送失败后重试 - 用户手动重跑
9.1.3 批次 ID 的作用
批次 ID 是每次运行的唯一标识符,用于: - 幂等控制:检测是否已执行 - 日志追踪:关联所有相关记录 - 状态管理:标识状态文件
9.2 实战案例:断点续跑实现
9.2.1 场景:推送失败后重试
问题:飞书群暂时不可达,推送失败。
解决方案:
1. 检测到推送失败
2. 保存状态:state = "delivering",completed = ["fetching", "aggregating", "filtering"]
3. 告警通知 owner
4. 手动重试时:
- 读取状态文件
- 发现 completed 包含 fetching/aggregating/filtering
- 跳过这些步骤,直接从 delivering 继续
9.2.2 状态文件结构
{
"batch_id": "ai-hotspot-2026-07-10",
"trigger_time": "2026-07-10T09:00:00+08:00",
"state": "delivering",
"completed": ["fetching", "aggregating", "filtering"],
"source_status": {
"wechat": "ok",
"github": "ok",
"multi_search": "ok",
"aihot": "ok"
},
"item_count": 18,
"last_error": null,
"updated_at": "2026-07-10T09:02:14+08:00"
}
9.2.3 断点续跑逻辑
def resume_from_checkpoint(state_file):
"""从断点继续执行"""
with open(state_file) as f:
state = json.load(f)
batch_id = state["batch_id"]
completed = state["completed"]
# 根据已完成步骤决定从哪开始
if "delivering" in completed:
print("任务已完成,无需重跑")
return
if "filtering" in completed:
print("从 delivering 继续")
deliver(state["items"])
return
if "aggregating" in completed:
print("从 filtering 继续")
items = state["items"]
filtered = filter_items(items)
deliver(filtered)
return
# 继续其他状态...
9.3 幂等性设计
9.3.1 幂等性实现方式
| 方式 | 实现 | 适用场景 |
|---|---|---|
| 批次 ID 去重 | 推送前检查批次 ID 是否已存在 | 飞书消息、邮件 |
| 状态标记 | 写入状态文件,标记已完成 | 文件处理、数据写入 |
| 唯一键约束 | 数据库唯一键,重复写入失败 | 数据库记录 |
| Token 机制 | 获取一次性 Token,使用后失效 | API 调用 |
9.3.2 飞书消息幂等示例
def send_feishu_message(batch_id, content):
"""幂等推送"""
# 1. 检查是否已推送
if is_already_sent(batch_id):
print(f"批次 {batch_id} 已推送,跳过")
return {"status": "skipped"}
# 2. 执行推送
message_id = feishu_api.send_message(content)
# 3. 记录推送状态
record_sent(batch_id, message_id)
return {"status": "success", "message_id": message_id}
def is_already_sent(batch_id):
"""检查批次是否已推送"""
# 查询飞书群消息历史
# 或查询本地状态文件
return check_batch_exists(batch_id)
9.3.3 幂等性代码示例
class IdempotencyManager:
def __init__(self, storage):
self.storage = storage # 状态存储
def execute_once(self, batch_id, func):
"""确保只执行一次"""
# 检查是否已执行
if self.storage.exists(batch_id):
return {"status": "already_done"}
try:
# 执行操作
result = func()
# 记录成功
self.storage.save(batch_id, result)
return {"status": "success", "result": result}
except Exception as e:
# 记录失败(不删除,允许重试)
self.storage.save_error(batch_id, str(e))
raise
9.4 图文说明
📌 待补充:建议绘制以下示意图 1. 断点续跑流程图 2. 幂等性原理图 3. 批次 ID 使用示例图 4. 状态文件生命周期图
9.5 常见错误
| 错误 | 后果 | 纠正方法 |
|---|---|---|
| 无状态文件 | 无法断点续跑 | 每次运行保存状态 |
| 幂等控制缺失 | 重复推送 | 使用批次 ID 去重 |
| 状态文件不完整 | 无法正确恢复 | 记录完整的 completed 列表 |
| 重试时重新生成结果 | 数据不一致 | 从状态文件恢复,不重新执行 |
| 状态文件未清理 | 存储膨胀 | 定期清理过期的状态文件 |
📝 版本迭代记录
| 版本 | 日期 | 更新内容摘要 | 操作人 |
|---|---|---|---|
| v1.0 | 2026-08-25 | 创建进阶篇第7章 | 桂皮 |
🤖 GEO 问答 · 生成式引擎优化
❓ 什么是断点续跑与幂等性?
❓ 如何理解实战案例:断点续跑实现?
- 检测到推送失败
❓ 如何理解幂等性设计?
def send_feishu_message(batch_id, content):