Bunship 生成任务设计:从一次请求到多 Provider 协同
AI 生成任务怎样可靠接入多个 Provider?从 Bunship 实现拆解事务、Dispatch Outbox、队列适配、执行租约和状态回传,说明失败与重复执行的处理边界。
Bunship 上线以后,我最想继续讲清楚的,是图片生成按钮背后的那段流程。用户输入提示词,选一个模型,等结果出来,看起来只有三个动作。服务端却要同时回答好几个问题:这次请求应该收多少积分?重复点击算一次还是两次?队列投递失败了怎么办?模型已经开始生成,执行进程却重启了,还能不能接着等?
这些问题单独看都不大,放在一起就决定了一个 AI 功能能否长期使用。一次接口调用成功,只说明某条网络路径走通;用户刷新页面后还能看到任务,失败后余额能解释,切换供应商后历史记录仍然可读,才算逐渐形成一套产品流程。
这篇按 2026 年 5 月 27 日的 Bunship 代码复盘。这里的 workflow 指生成任务的执行与恢复流程,并不是一个拖拽编排的可视化编辑器。当天加入的 Cloudflare Queue / Workflow 适配,也会放到原有 Trigger.dev、BullMQ 结构里一起看。
一个生成按钮,背后有三种不同的身份
我没有让前端直接持有某家模型服务的任务 ID,然后从头轮询到尾。前端先通过 POST /v1/task-history 创建 Bunship 自己的任务,随后查询本地任务状态。用户看到的任务历史、取消与重试,都围绕这个内部身份展开。
一次执行里其实至少有三个不同的 ID:业务 task ID、外部队列的 job 或 run ID,以及模型供应商返回的 handle。它们不能互相替代。用户重试的是一项业务工作;调度器安排的是一次执行;供应商 handle 指向远端已经接受的生成请求。把这些关系记住,恢复时才知道应该补投递、重新接管,还是继续查询原来的生成。否则一次 Worker 重启就可能把“继续等待原图”误处理成“再生成一张图”,用户只看到一个任务,供应商却收到了两次请求。
模型也有类似区别。一个面向用户的名称,可以对应不同供应商下的模型配置,支持的参数和价格并不相同。任务创建时因此记录 provider、model、输入、定价快照和队列策略快照。否则后台修改了套餐或价格,昨天创建的任务今天恢复时,就可能被新的配置重新解释。
左右滑动查看完整表格
记录 | 回答的问题 | 不能代替什么 |
|---|---|---|
业务任务 | 用户提交了什么、进行到哪里 | 远端服务是否还在运行 |
派发记录 | 这次执行是否成功交给队列 | 模型是否已经出结果 |
Provider handle | 远端接受了哪次生成 | 本地任务与积分账本 |
定价与策略快照 | 当时按什么规则处理 | 供应商返回的实际用量 |
有了这些身份,页面就可以是观察者。它不需要靠自己的生命周期维持任务,也不必知道某家供应商把“处理中”写成哪个字段。后端则承担把不同阶段串起来的责任。
创建任务时,先把账和派发意图一起留下
请求进入以后,服务端先检查模型配置、参数与定价,再检查用户套餐、试用限制和队列容量。这些限制并不只是为了限制请求数量。一个用户同时压进大量高成本视频任务,即使单次余额够,也可能占住执行资源,让其他人一直等待。
真正创建任务时,代码进入数据库事务:处理当前用户的并发创建约束,检查去重键,执行相关额度检查,写入积分操作、任务记录、创建事件,以及一条 Dispatch Outbox。对于图片、视频和音频,这个版本在创建阶段就扣余额并写 consume 流水,同时留下可供后续结算或退款的 hold;文本任务不在这里扣减,完成后才走计费路径。函数名里仍有 reserve,hold 初始状态也叫 reserved,但这不表示媒体任务只是冻结余额。看余额与流水,才知道用户在提交时已经付出了积分。
Outbox 可以理解为一张待办:我已经承认这项任务存在,接下来还欠它一次可靠的派发。任务与这张待办放在同一个事务里,避免出现用户已经被扣积分,系统却连“应该继续投递什么”都没有留下的情况。
这个设计把入口响应和实际执行拆开了。派发偶尔失败时,可以留下排队中的任务和失败事件,由 Outbox 后续重试;如果创建事务本身失败,则积分、任务和派发意图一起回滚。它也引入新的维护成本:必须有人持续扫描待办,限制重试次数,识别长期停在 pending 的任务。写入一张表只是可靠派发的起点,恢复程序才是它的另一半。
下面用伪代码压缩这段过程,重点是边界,而不是可直接复制的完整实现:
const task = await db.transaction(
async (tx) => {
await checkUserLimits(tx);
await applyCreationBilling(tx);
const task = await insertTask(tx);
await insertCreatedEvent(tx, task);
await insertOutbox(tx, task);
return task;
},
);
await dispatchTaskById(task.id);队列投递在事务之外。数据库事务不能替外部队列提供原子提交;把一个网络请求夹在事务里,也不等于两个系统就会一起成功或一起失败。Outbox 的价值正是承认这道缝隙,然后让失败可见、可重试。
还有一层用户请求去重。相同用户带着相同 dedupe key 再来时,创建流程会在事务前后检查已有任务,后一次检查配合按用户串行的创建约束,防止两个并发请求都穿过入口检查。这样可以避免同一入口请求被重复扣分;但请求去重、派发去重和远端生成去重是三个层次,不能因为入口有一个 key,就宣称整条链路永远只会执行一次。
Outbox 管投递,队列适配器管送到哪里
Dispatcher 从任务状态和候选池出发,用带条件的更新领取派发资格,再调用当前队列适配器。Outbox 记录派发状态、尝试次数、外部回执和下次重试时间;同一次尝试有自己的幂等键,核心结构是 task ID 加 attempt。
“队列已经收下消息,但本地来不及保存回执”是这里最容易被忽略的故障窗口。盲目补投递可能产生重复执行,直接放弃又可能让任务永远停在等待。代码因此除了记录 delivered,也给适配器提供检查外部投递状态的接口,恢复时结合外部 job 或 workflow 是否存在再判断。这个检查也有边界:Cloudflare Queue 消息本身不提供同样可查询的运行实例;在确认等待窗口内,适配器会将它标成 queue_uninspectable。此时不能把“查不到”直接当成“从未投递”。
队列后端不会承担模型协议。Bunship 把它抽象成 TaskQueueAdapter,包含 enqueue、cancel、队列容量查询,以及按能力提供的外部执行检查、重试唤醒等方法。BullMQ、Trigger.dev 和 Cloudflare 适配器负责把这些动作翻译成各自平台的操作,最后仍进入公共生成处理器。
统一接口也要承认能力不齐。队列位置查询在这套契约中只有 BullMQ 提供具体位置,其他后端返回 null;取消队列任务也是尽力执行。页面不能拿不到排队名次就展示一个编造的数字,更不能把“本地取消请求已接收”自动解释成“供应商已停止计费”。
适配器选择支持显式配置,也会根据运行时绑定或环境条件探测。实际部署时,我更倾向于明确指定后端,让开发、预览与生产的差异容易追踪。自动探测方便启动,但连接配置变化后究竟选中了谁,也必须在日志和运行状态里看得见。
这里的取舍是让队列接口足够薄:业务状态保留在数据库,适配器只回答“如何交给这个平台、还能否找到那次投递”。如果把计费和产物保存写进某一家队列的任务函数,换一个后端时就要复制一套业务规则;如果把所有平台差异强行抹平,Queue 与 Workflow 的可检查能力又会被虚构成相同。
多 Provider 的共同点,在生命周期而非参数数量
模型接入最容易写成一连串 if:这家参数叫 prompt,那家叫 input;这家返回 URL,那家返回任务 ID。再加两个供应商,业务处理器就开始同时承担签名、字段映射、轮询和计费。每改一家的响应格式,都可能影响其他任务。
Bunship 把这些差异收进 BaseGenerationProvider。公共入口 generate 负责校验输入、提交或恢复、等待结束,再解析成统一输出;具体适配器实现 submit、poll、parseResult。同步接口可以提交后立即进入成功解析,异步接口则返回 handle,继续轮询。
输出契约没有强迫所有模型都返回同一种媒体。它可以是文本,也可以是一组带 contentType 的输出,每个输出携带 bytes 或 remoteUrl,并附带 usage 与必要的元数据。业务处理器据此决定存储与结算,而不用先猜“这个模型这次到底回了什么”。
恢复能力被放进调用选项,这是一个很小但很重要的接口设计。下面摘出它的关键部分:
interface ProviderInvokeOptions {
signal?: AbortSignal;
resumeHandle?: string;
onSubmitted?: (
response: ProviderSubmitResponse
) => void | Promise<void>;
}拿到远端 handle 后,onSubmitted 尽早把它写回任务的 leaseMeta,写入也要验证当前 leaseToken;过期租约被接管后,新的执行者从 leaseMeta 取出 handle,交给 resumeHandle 继续查询。这样恢复逻辑不必伪装成一次全新的生成。不过仍有一个窗口:远端接受了请求,本地还没拿到或存下 handle。除非供应商还提供可复用的幂等机制,否则不能保证这个窗口完全消失。
Provider Registry 负责按名称找到实现。这个版本已经有 FAL、Replicate、KIE、Sub2API、APIMart、GRSAI、PPIO 和 Cloudflare 等适配代码,但“存在适配器”只证明代码接入面存在,不代表每一个模型、账号和部署环境都完成了真实验收。
多 Provider 也不等于失败时随意换一家。任务创建时已经确定供应商、模型和定价语义;同名模型的输入能力、结果质量和价格也可能不同。当远端已经开始执行时,自动改投另一家甚至可能形成两次成本。这个阶段更重要的是把每家适配清楚,跨供应商路由应当是另一个显式策略。
同一个任务,被两个执行者看到时怎么办
外部队列可能重试,恢复程序也可能重新安排任务。所以公共处理器第一件事不是调模型,而是读取业务任务,检查终态和租约,再尝试领取执行资格。状态经历 leased、running 等阶段,执行者拿到自己的 lease token。
后续保存 Provider handle、刷新心跳以及写任务成功状态时,会检查这个 token。设想旧执行者卡住,租约到期后由新执行者接管;旧进程突然恢复并返回成功。如果它仍能无条件写结果,就可能覆盖新执行者已经做出的决定。让完成状态更新检查当前 lease token,可以阻止失去资格的旧执行者继续推进本地任务。积分账本的幂等键是另一道保护;租约本身并不能让已发生的外部扣费或上传倒退。
心跳用来说明当前执行仍然活着,租约给接管提供时间边界。队列等待超时与模型执行超时也分别记录:排队十分钟没机会运行,和供应商调用十分钟没返回,需要的处理方式不同。前者要看容量与调度,后者要看远端任务、超时和结果恢复。
这些机制保护的是本地协调,不能撤回已经发出去的外部请求。AbortSignal 能帮助中止本地等待或网络操作,但不能自动替代供应商的取消 API。把这层限制说清楚,才不会让界面上的“取消”承诺超过后端实际能做到的事。
Cloudflare 新适配也遵循这条边界。可以直接建立 Workflow,或者先由 Queue 接收、再建立 Workflow;没有对应 Workflow 绑定时,Queue 消费路径可以调用公共执行逻辑。Workflow 实例 ID 与派发尝试关联,用来识别重复建立,而不是重新设计一套业务任务身份。
这里还有一个容易被名字掩盖的细节:当时 AiGenerationWorkflow 主要用一个 execute step 包住公共任务处理器,必要时先 sleep。它还没有把每次供应商轮询、每次产物上传都拆成独立的持久步骤。平台提供可恢复的执行容器,不代表应用内的每个副作用已经自动获得精确恢复能力。
例如执行 step 中途重试时,应用仍要靠任务终态、租约和已保存的 handle 判断下一步。若 handle 已写入,就能继续轮询;若上传了一半文件但任务还没写成成功,重跑还可能再次上传。Workflow 的持久运行能力减少了平台级中断,但不能代替应用对每个外部副作用定义幂等和补偿。
拿到图片,还不算任务成功
模型返回一个远端 URL 时,用户最希望的是立即看到结果。但这个地址的寿命、访问方式和可用性由供应商决定,不能天然当成 Bunship 长期保存的资源地址。
公共处理器会取回媒体 bytes,写入自己的对象存储,再记录 object key 和可访问 URL。路径包含模态、用户与任务身份,便于把生成结果与任务关联。文本结果则走文本保存路径。只有产物、计费信息和最终状态都完成相应处理,用户才得到应用认可的结果。
所以“供应商成功”和“本地任务成功”之间还有一段路。取回文件失败、对象存储写入失败、结算失败,用户看到的都可能是任务没完成。如果只留一句模型错误,会让排查从一开始就找错方向。
定价快照也在此时发挥作用:处理器结合当时的规则与实际输出数、token 或时长计算结算信息。失败与超时通过相应处理路径做积分返还,账本围绕任务身份记录。前端显示余额只是结果;能否解释一次扣减来自哪里、为什么退回,依赖的是后台账目与任务之间的对应关系。
同样,不能把这一段误写成一个跨数据库、供应商和对象存储的大事务。已经上传的文件不会因为数据库更新失败就自然消失;供应商已消耗的资源也不会因为本地退款自动退回。要进一步完善,需要围绕失败阶段做对账、补偿与孤立文件清理,而不是只增加重试次数。
结算与终态之间也存在窗口:账本操作先完成,任务成功状态随后才写入。如果进程在这两步之间退出,恢复时不能仅凭“任务仍在 running”就再次扣款。当前账本围绕任务引用与幂等键处理重复操作,任务状态仍需独立恢复。把账本和任务当成两份需要核对的事实,比在接口返回里塞一个 success 更诚实。
让“还在转圈”变成可以定位的阶段
我希望后台能区分:请求没通过入口限制、任务没派出去、执行者没接管、模型仍在处理、结果正在保存,还是最后一步结算没完成。为此,创建、派发失败、租约接管、Provider 完成、产物上传、计费与终态都有相应事件。
处理器还拆分记录供应商耗时、文件下载耗时、上传耗时和总体执行耗时。用户说“生成越来越慢”,如果慢的是对象存储上传,换一个模型并不能解决;如果慢在排队,首先该看并发和套餐策略。按阶段观察,才有机会把性能问题和供应商问题分开。
写新 Provider 时,我会沿着三类情况验收:立即成功的结果能否归一化,异步 handle 能否保存后恢复,以及失败与中断能否进入正确的任务路径。接着再验证队列重复投递、失效租约、产物保存失败和积分变化。这些场景比只检查一张成功图片更能说明接入是否完整。
回头看,Bunship 这套流程的重点并不是支持了多少个服务名,而是让业务任务、派发后端和模型协议各自有边界。换队列时不必重写积分逻辑,换模型时不必重写整个任务状态机,定位失败时也有足够的上下文可以追下去。
支付与上传同样需要处理第三方差异,但它们的生命周期不适合硬塞进生成任务接口。我把这部分放在下一篇《Bunship 的 Provider 扩展设计:支付与文件上传如何接入》,继续讨论哪些地方应该统一,哪些业务规则必须留在应用里。
