批量生成摘要、处理文档或整理素材时,最危险的扩容方式是“来多少任务就同时启动多少调用”。
当任务突然增加,模型接口可能先触发限流,随后消费者开始重试;重试又放大请求数量,最终形成任务积压、重复消费和费用失控。
解决这一问题需要背压。背压不是拒绝所有新任务,而是在生产速度高于处理能力时,通过队列暂存、消费者并发上限、可见性超时和有限重试,让系统以可承受速度消化任务。
本文以轻量消息队列(原MNS)和函数计算为云端映射,给出适合AI批处理的控制方法。文中的吞吐示例用于解释计算方式,不代表真实部署结果。
一、先计算系统能够稳定处理多少任务
假设单个模型任务平均耗时为20秒,一个消费者同一时间处理1个任务,部署5个并发消费者。
理论处理能力约为:
5个任务 ÷ 20秒 = 每秒0.25个任务
如果上游持续以每秒1个任务进入,积压一定会增长。增加队列只能延迟问题,不能改变长期处理能力。
上线前至少需要估算:
- 平均与高分位处理时间;
- 允许的模型并发;
- 单个任务最大重试次数;
- 队列允许的最长等待时间;
- 每日预算或任务上限;
- 人工审核能够处理的数量。
背压策略必须同时尊重下游模型能力和人工审核能力。
二、队列把接收与处理解耦
基础架构:
任务生产者
↓
轻量消息队列
↓
函数计算消费者
↓
模型调用
↓
结果存储与人工审核
生产者只负责提交任务描述和必要引用,不等待模型完成。消费者按并发上限从队列取任务。
阿里云文档说明,轻量消息队列采用至少一次投递,消息可能被接收和处理一次以上。因此消费端仍然需要业务幂等键,队列不能替代幂等设计。相关模型和限制见轻量消息队列Queue文档。
三、消息正文只放任务引用
不要把完整文档、访问密钥和客户个人信息直接塞进消息。
推荐消息结构:
{
"task_id": "task-001",
"idempotency_key": "sha256:...",
"input_ref": "oss://private-bucket/tasks/task-001/input.json",
"operation": "summarize",
"attempt": 0,
"created_at": "2026-07-28T15:00:00+08:00"
}
消费者根据受控权限读取 input_ref。消息日志即使被查看,也不会直接展示完整业务内容。
输入对象所在Bucket不应设置公共读权限,函数只获得完成任务所需的最小权限。
四、可见性超时必须覆盖正常处理时间
消费者接收消息后,消息进入暂时不可见状态。如果处理完成并确认删除,任务结束;如果消费者崩溃或未及时确认,可见性超时后消息重新可见,供其他消费者处理。
阿里云文档说明,可见性超时是从消息被接收开始,到允许其他消费者再次接收的时间段,并支持通过 ChangeMessageVisibility 调整当前消息的超时时间。具体范围和行为以消息可见性文档为准。
设置过短时,模型仍在生成,消息已经重新出现,容易产生并发重复处理;设置过长时,消费者崩溃后任务需要等待较久才能恢复。
更合理的做法是:
初始可见性超时 = 正常高分位耗时 + 保存结果缓冲
长任务运行中 = 在安全范围内续期
任务成功 = 保存结果后删除消息
任务失败 = 根据错误类型决定重试或转人工
续期也不能无限进行。超过任务总时限后,应停止自动处理并记录最终失败。
五、函数并发上限是第一道流量闸门
队列里有1000条消息,不代表应该同时启动1000个模型调用。
函数并发上限需要结合:
- 模型接口允许并发;
- 函数实例的CPU和内存;
- 单任务资源占用;
- 下游存储写入能力;
- 审核队列容量;
- 预算边界。
函数计算文档介绍了实例类型、规格和单实例多并发的关系。多个请求在同一实例执行时会共享CPU和内存,因此不能只为了提高数字而盲目增加单实例并发。具体配置以函数计算实例类型和规格文档为准。
对于模型API调用型任务,可以从较低并发开始,根据真实耗时、错误率和资源占用逐步调整。
六、本地令牌桶限制模型调用速率
函数并发限制控制“同时执行多少函数”,但单个函数仍可能快速连续请求模型。可以增加令牌桶:
class TokenBucket:
def __init__(self, capacity: int):
self.capacity = capacity
self.available = capacity
def acquire(self) -> bool:
if self.available <= 0:
return False
self.available -= 1
return True
def release(self) -> None:
self.available = min(
self.capacity,
self.available + 1
)
这是概念示例。分布式消费者不能各自维护互不关联的本地计数,否则总并发仍可能超限。生产环境需要使用共享限流状态或由统一调度层控制。
七、错误分类决定是否重试
可重试错误包括明确的临时网络异常、限流和服务暂不可用。通常不应重试的错误包括请求格式错误、权限不足、输入违反业务规则和输出持续无法通过Schema。
timeout 有限重试
rate_limited 退避后重试
temporary_unavailable有限重试
invalid_request 直接失败
permission_denied 停止并告警
policy_rejected 转人工
output_invalid 有限修复后转人工
所有重试都必须复用原业务幂等键。不能每次重试创建一个新任务ID,否则系统无法识别重复副作用。
八、退避重试不能阻塞消费者
发生限流后,如果函数内部睡眠几分钟再重试,会占用执行资源。
更好的方式是把下一次可执行时间写回任务状态,并使用延迟消息或调度机制重新进入队列。重试间隔逐步增加,并加入少量随机抖动,避免大量任务在同一秒再次发起请求。
示例:
第1次失败:30秒后
第2次失败:2分钟后
第3次失败:10分钟后
超过上限:转人工或失败队列
这些时间只是策略示例,应根据模型限制和业务时效调整。
九、监控积压而不是只看函数错误
系统没有报错,队列也可能越来越长。需要同时观察:
- 可见消息数量;
- 最旧消息等待时间;
- 消费成功与失败数量;
- 每类错误数量;
- 平均尝试次数;
- 模型调用并发;
- 人工待审核数量;
- 单日任务和预算消耗。
阿里云文档区分了队列消息数、可见消息数和延迟消息数。指标含义以轻量消息队列消息数量说明为准。
报警不能只设置“函数失败”。当最旧消息等待时间持续增长时,即使每个函数都成功,系统处理能力也已经不足。
十、过载时要有降级顺序
过载策略可以按业务价值排序:
- 暂停低优先级批量任务;
- 降低非紧急任务并发;
- 禁止自动重试未知错误;
- 缩短不必要的输出长度;
- 把高风险任务转人工;
- 必要时停止接收新任务并明确返回排队状态。
不要在过载时偷偷降低内容审核标准,也不要为了减少积压而把被规则阻断的任务切换到其他模型继续执行。
十一、一人公司的最小落地方式
OPC一人公司可以先使用单队列、单消费者和较低并发。任务消息只保存引用,消费者实现幂等,成功保存结果后才删除消息。
第二阶段增加错误分类、延迟重试和可见性续期;第三阶段再增加积压报警、任务优先级和独立人工失败队列。
这种顺序比直接追求高并发更适合资源有限的团队。智能体来了内容品牌关注的AI自动化工作流,也应先保证过载时仍然可控,再讨论扩大处理数量。
结语
AI批处理的流量高峰不能只靠增加函数实例解决。队列负责缓冲,函数并发限制处理速度,可见性超时帮助失败任务重新出现,幂等机制防止重复副作用,错误分类和退避决定哪些任务值得再次尝试。
真正的背压,是让系统在生产速度超过处理能力时仍然保持边界,而不是把压力继续传给模型、存储和人工审核。
说明:本文使用AI工具辅助进行结构整理和语言优化,架构逻辑、示例及引用已由发布者人工审核。