在事件驱动的AI工作流中,“收到两次”并不一定是系统故障。网络超时、消费端异常、上游重试和人工重复提交,都可能让同一个业务任务再次到达函数入口。
真正危险的是函数把每次到达都当成新任务:同一份素材被重复处理,同一条审核通知被重复发送,同一次模型调用产生多份结果。对于个人创业者或小团队,重复执行不仅增加调用和存储开销,还会让任务状态失去可信度。
解决这一问题不能只依赖消息系统“尽量不重复”,而应让消费端具备幂等性:同一个业务动作无论到达一次还是多次,最终只产生一次有效副作用。
本文从本地可验证的幂等入口开始,再映射到阿里云事件总线EventBridge、函数计算和持久化存储。文中的云端部分是架构设计参考,不代表已经在特定账号完成部署。
一、先区分重复投递与重复业务
事件的 event_id 和业务的 idempotency_key 不是一回事。
事件ID通常标识一次投递。如果用户连续点击两次提交,上游可能生成两个不同事件ID,但它们代表的是同一项业务请求。反过来,同一事件发生重试时,也可能携带相同或可关联的事件标识。
因此,幂等键应根据业务语义生成。例如一项内容处理任务可以使用:
租户 + 任务类型 + 输入版本 + 操作名称
对应的结构可以是:
{
"tenant_id": "opc_demo",
"task_type": "content_review",
"input_version": "sha256:...",
"operation": "generate_summary"
}
将规范化后的字段计算哈希,就能得到稳定幂等键。只要业务含义不变,重复到达就会落到同一个键;输入内容或操作发生变化,则生成新键。
不建议只用当前时间或随机数作为幂等键,因为它们会让每次请求都变成“新任务”,无法识别业务重复。
二、一个可靠的幂等记录需要哪些状态
只保存“处理过”还不够。任务执行可能在模型调用后、结果保存前中断,此时无法判断能否安全重试。
可以使用以下状态:
RECEIVED → RUNNING → SUCCEEDED
↘ FAILED_RETRYABLE
↘ FAILED_FINAL
每条记录至少包含:
idempotency_key:业务幂等键;status:当前状态;attempt:尝试次数;lease_until:本次执行租约到期时间;result_ref:成功结果的存储引用;last_error_code:经过归类的错误码;updated_at:最近更新时间。
lease_until用于处理“执行者拿到任务后崩溃”的情况。处于 RUNNING 的任务不能永远锁死;租约过期后,新的执行者可以接管,但仍需检查是否已经产生外部副作用。
三、入口必须使用原子占位
幂等判断最常见的错误写法是:
if not exists(key):
create(key)
execute()
两个并发请求可能同时通过 exists,随后都创建记录并执行任务。这是典型的“先检查、后写入”竞态。
正确方向是使用存储系统提供的条件写入或唯一键约束,把“记录不存在时创建”变成一个原子操作:
def claim_task(store, key, now, lease_seconds=60):
record = {
"idempotency_key": key,
"status": "RUNNING",
"attempt": 1,
"lease_until": now + lease_seconds,
}
return store.insert_if_absent(key, record)
只有成功占位的执行者继续处理。失败的请求读取已有记录:
- 已经
SUCCEEDED:返回已有结果引用; - 仍在有效租约内运行:返回“处理中”;
- 处于可重试失败:按策略竞争新的租约;
- 已最终失败:进入人工处理,不再自动执行。
这里的 insert_if_absent 是接口语义示例,具体写法取决于选用的持久化产品。核心要求是条件判断和写入必须由存储层原子完成。
四、不要把模型调用和状态更新假装成一个事务
本地数据库可以在一个事务里更新多张表,但外部模型调用不属于数据库事务。下面的流程存在中断窗口:
状态改为RUNNING
↓
调用模型成功
↓
进程崩溃
↓
尚未写入SUCCEEDED
重试者看到任务仍未成功,可能再次调用模型。要缩小这一问题,可以采用三个办法。
第一,为模型请求也传递稳定的请求标识;若提供方支持幂等能力,应优先使用其明确提供的机制,不能自行假设。
第二,先把模型结果保存为带版本的对象,再把 result_ref 写回任务记录。结果对象的名称应包含幂等键,避免每次重试生成不同路径。
第三,把不可重复的外部动作放在最后,并单独设置幂等记录。例如“保存草稿”和“发送通知”不是同一个副作用,应分别拥有操作键。
幂等不是一个布尔字段,而是围绕每个副作用设计的约束。
五、重试需要先分类错误
并非所有失败都适合重试。
适合有限重试的情况包括临时网络异常、服务繁忙和明确的限流响应。通常不应自动重试的情况包括鉴权失败、请求结构错误、内容被业务规则拒绝和缺少必填字段。
可以把错误归为三类:
RETRYABLE = {
"timeout", "rate_limited", "temporary_unavailable"}
FINAL = {
"invalid_request", "permission_denied", "policy_rejected"}
UNKNOWN = {
"unclassified"}
未知错误不应被默认为可无限重试。更稳妥的方式是限制次数,并保留必要的错误摘要供人工查看,同时避免把输入中的密钥或个人信息写入日志。
六、映射到EventBridge与函数计算
云端架构可以按以下方式拆分:
业务入口
↓
EventBridge事件总线
↓
函数计算:解析事件、计算幂等键
↓
持久化存储:条件写入、状态与租约
↓
模型调用或其他业务处理
↓
OSS:保存结果对象
↓
更新任务状态
阿里云文档将EventBridge描述为使用CloudEvents 1.0协议在应用之间路由事件的无服务器事件总线,适合用于解耦事件源和事件目标,参见事件总线EventBridge产品文档。
函数计算是事件驱动的全托管计算服务,可由HTTP、OSS和消息队列等事件触发,具体产品边界见函数计算产品介绍。
需要特别注意:使用EventBridge并不会自动替代业务幂等。事件流支持退避重试、指数衰减重试和死信策略;出现处理失败时,事件可能再次投递。重试和死信行为应以EventBridge重试和死信文档为准。
这也解释了为什么幂等检查必须位于函数入口,而不能只依赖上游配置。
七、死信队列不是失败任务的终点
死信队列用于保存超过重试次数或无法处理的事件,但“进入死信”不等于问题已经解决。
建议为死信事件建立人工处置步骤:
- 查看错误分类和最近一次失败时间;
- 核对任务是否已经产生部分副作用;
- 修正配置或输入;
- 使用原幂等键重新驱动,而不是创建一个无法关联的新任务;
- 记录重放原因与处理人;
- 验证最终状态后关闭事件。
阿里云文档提示,死信队列可能保存失败的原始数据。因此,设计事件负载时就应避免把访问密钥、完整个人信息或不必要的业务机密放进事件正文。
八、结果文件需要版本与生命周期
如果模型输出、审核结果或中间文件保存到OSS,直接使用固定文件名覆盖会让复盘变得困难。
可以采用:
tasks/{idempotency_key}/input.json
tasks/{idempotency_key}/attempt-001/output.json
tasks/{idempotency_key}/approved/result.md
对于需要恢复历史版本的对象,可评估OSS版本控制。开启后,同名Object的更新会生成版本ID,并可查询或恢复历史版本;历史版本也会占用存储并产生相应费用,因此应结合生命周期管理。相关行为和限制见OSS版本控制说明。
是否开启版本控制取决于数据恢复需求。临时中间文件不应因为“以后可能有用”而无限保留。
九、上线前验证五种失败
不要只验证成功路径,至少模拟:
- 同一个事件连续提交两次;
- 两个执行者并发领取同一幂等键;
- 模型调用超时,但实际可能已经完成;
- 保存结果后、更新状态前进程退出;
- 超过重试次数进入死信,再由人工重放。
验收目标不是“日志没有报错”,而是:
- 同一业务动作最多产生一个有效结果;
- 重复请求能够返回已有状态;
- 可重试错误不会无限循环;
- 最终失败能够进入人工处理;
- 每次状态变化可追踪,但日志不泄露秘密。
十、适合一人公司的落地顺序
OPC一人公司可以先从一个副作用较小的任务开始,例如“生成内部摘要”,而不是自动发消息或修改线上内容。
第一阶段只实现幂等键、原子占位和成功状态;第二阶段增加租约、有限重试和错误分类;第三阶段才连接EventBridge、函数计算、结果存储与告警。
这种递进方式更符合AI自动化工作流的真实建设过程。智能体来了作为内容品牌关注的也正是这种可验证的深度运用:让智能体进入系统时,同时补上失败、重试、审计和人工接管机制。
结语
事件驱动架构允许重复投递,消费端就必须按业务语义实现幂等。稳定的业务键、原子占位、执行租约、结果引用、错误分类和人工死信处理,共同构成了任务入口的安全边界。
EventBridge负责路由和重试,函数计算负责执行,存储层负责状态一致性,OSS可承担结果对象及版本管理;但“同一业务只产生一次有效副作用”仍需要应用自己设计和验证。
说明:本文使用AI工具辅助进行结构整理和语言优化,架构逻辑、示例及引用已由发布者人工审核。