耗时的 AI 流程怎么异步执行?进度怎么推给前端?
简化版
接口只负责建一条任务记录(状态「运行中」或「待执行」)并把任务交给后台,立刻返回任务 ID;后台任务逐步执行,每到一个阶段就把「阶段描述 + 进度百分比」写进任务记录;前端拿着任务 ID 看进度,简单场景用定时轮询,要实时看每一步就用 SSE(服务端推送事件流),连接在任务进入结束状态后断开。后台执行要处理好四件事:同一个任务不能被调度两次;同时跑的任务数要限流;任务结束(包括失败)时状态一定要落库;异步线程拿不到请求线程里的上下文(当前用户、类加载器等),要显式传进去并在结束时清理。
详细版
POST /task/start → 校验 → INSERT task(status=RUNNING, progress=0) → 交给后台 → 返回 task_id
后台:阶段 1 → UPDATE progress=20, stage=「音频处理完成」
阶段 2 → UPDATE progress=55, stage=「正在转写第 2/3 段」
完成 → UPDATE status=SUCCEEDED, progress=100 (失败 → status=FAILED, error=…)
前端:轮询 GET /task/{id}(每 2–3 秒) 或 SSE GET /task/events/{id}(有变化就推)
| 方式 | 原理 | 优点 | 缺点 |
|---|---|---|---|
| 轮询 | 前端定时请求任务状态 | 实现最简单,无长连接 | 有延迟;任务多时请求量大 |
| SSE | 服务端保持一个 HTTP 响应,持续写出事件 | 单向实时、基于普通 HTTP、浏览器原生支持 | 服务端要维持连接;经过反向代理要关闭缓冲 |
| WebSocket | 双向长连接 | 双向实时 | 协议和运维更重,进度推送用不上双向 |
后台执行四件事:调度去重、并发限流、终态落库、上下文传递与清理。
完整版教学
一、为什么不能让请求等着跑完
一条 AI 流程往往包含多次模型调用,时间很容易超过 HTTP 链路上各层的超时:
自检 2 轮 = 审查 + 重写 + 审查 + 重写 = 4 次调用,每次约 8 秒(示意)→ 约 32 秒
常见网关或反向代理的读超时在 30–60 秒量级
→ 请求可能在任务跑完之前就被中间层掐断,前端报错,但后端还在继续算
更麻烦的是,请求被掐断后前端不知道任务到底成没成,用户重试就会再发起一次,同一件事做两遍。所以耗时流程应当改成「提交任务 → 立即返回 → 查进度」的异步模式,任务的真实状态以数据库里的记录为准。
二、后台任务的基本骨架
不管用线程池、协程还是消息队列,后台任务都要满足三条:
| 要求 | 做法 | 不做的后果 |
|---|---|---|
| 同一任务不重复调度 | 维护「任务 ID → 运行中的任务」的注册表,已在跑就直接返回;或用「只有待执行才能被抢占」的条件更新 | 同一份数据被处理两遍,结果互相覆盖 |
| 注册表不无限增长 | 任务结束的回调里把自己从注册表摘掉 | 内存泄漏 |
| 终态一定落库 | 最外层捕获所有异常,写「失败 + 原因」 | 记录永远停在「运行中」 |
以 Python 协程为例,骨架只有几行:
def schedule(task_id):
running = _running.get(task_id)
if running and not running.done(): # 已在跑:不重复调度
return
t = asyncio.create_task(execute(task_id))
_running[task_id] = t
t.add_done_callback(lambda _: _running.pop(task_id, None)) # 结束即摘除
注意注册表只在当前进程内有效,进程重启后就空了,所以「是否在跑」最终还要以数据库状态为准,服务启动时要把停在「运行中」的记录重新接管。断点怎么存、重启后怎么接着跑,见「长流程怎么支持中断和恢复?断点存在哪?」。
三、并发限流:排队比一起上更快
模型接口通常有并发和速率限制,一次性放开所有任务反而会大量报错。用信号量限制同时执行的任务数,其余任务在门口排队:
10 个任务,每个约 60 秒,信号量 = 2
分 5 批完成,最后一个任务约在第 300 秒结束
排队中的任务状态仍是「待执行」,页面显示「等待执行」
不限流:10 个同时打到模型接口,触发限流后部分失败,还要重试
信号量要在事件循环里创建,并且只创建一次,比如第一次用到时再创建,保证它绑定在应用实际运行的那个事件循环上。排队的任务状态要和正在执行的区分开,用户才知道「不是卡住了,是在排队」。
记忆钩子:异步任务四件事——不重复、要限流、终态落库、上下文带进去。
四、进度怎么算才不骗人
进度由两部分组成:给人看的阶段描述(「正在转写第 2/3 段」)和给进度条用的百分比。百分比的常见算法:
按阶段分区间:准备 0–20,逐段处理 20–90(按段数均分),收尾 90–100
共 3 段:第 1 段开始 20,第 2 段开始 20 + 70/3 ≈ 43,第 3 段开始 20 + 140/3 ≈ 66
按步骤数:progress = min(95, 已开始步数 × 100 / 最多步数)
最多 5 步:第 3 步开始时 60;收尾时再单独写 100
封顶 95 的原因是:最多步数是按最坏情况算的,提前通过时步数用不完,如果不封顶、也不在收尾时置 100,进度条会停在一个奇怪的数字上;反过来如果中间就到了 100,用户会以为已经完成。进度只能单调增加,阶段描述要具体到「第几段、第几轮」,比一个转圈的图标有用得多。
五、轮询还是 SSE
两种方式可以同时用在不同页面:列表页用轮询,只有列表里存在「运行中」的行时才发请求;详情页用 SSE,看每一步的实时变化。
| 场景 | 推荐 | 理由 |
|---|---|---|
| 列表里有若干条任务在跑 | 轮询,2–3 秒一次,没有运行中的行就不发 | 实现简单,多条任务共用一次请求 |
| 打开一个任务看步骤时间线 | SSE | 每一步变化立刻可见,不用高频请求 |
| 需要前端向运行中的任务发指令 | 普通接口(中断、恢复) | SSE 是单向的,指令走独立接口即可 |
SSE 的报文格式很简单:每条以 data: 开头、以两个换行结尾。服务端的写法通常是一个循环:读一次任务记录,内容和上次推送的不同才推,进入结束状态就退出循环、关闭连接,否则等一秒再读。
六、SSE 落地的几个细节
响应类型: Content-Type: text/event-stream
禁止缓存: Cache-Control: no-cache
代理缓冲: 经过 Nginx 时关闭缓冲,否则报文攒一批才到浏览器
去重推送: 和上次推送内容相同就不推,省流量也省前端渲染
主动断开: 关闭弹窗、离开页面、发出中断指令时,前端要主动断开连接
前端读取时,一次读到的数据块可能只有半条报文,要用缓冲区拼接,按两个换行切分,最后不完整的那段留到下一块到达再处理。推送的内容可以是完整的任务数据(包含步骤列表),前端整体替换,比推增量更不容易出错。
七、异步线程的上下文陷阱
把任务交给后台线程后,很多「理所当然」的东西就没了:
| 丢失的东西 | 后果 | 做法 |
|---|---|---|
| 当前登录用户 | 日志、审计记录里没有操作人 | 发起时把用户 ID 作为参数传进去 |
| 线程级上下文(ThreadLocal) | 线程池复用时,上一个任务的值残留到下一个任务 | 用完在 finally 里清理 |
| 线程的上下文类加载器 | 在打包运行的环境里,部分依赖初始化失败 | 在后台线程开头设置成应用自己的类加载器 |
| 框架代理 | 同一个类里自己调用自己的异步方法,异步不生效 | 异步方法放在单独的类里,从外部调用 |
还要注意异常类型:后台代码常写 catch (Exception e) 来记失败,但类初始化失败这类问题抛出的是 Error,接不住,任务记录就会一直停在「运行中」。
八、常见误区与追问
- 误区:把 HTTP 超时调大就不用异步了。 链路上网关、代理各有超时,请求断开后后端还在算,用户重试会导致重复执行。
- 误区:后台任务不用限流,越多越快。 模型接口有并发和速率限制,一次性放开会大量失败。
- 误区:内存里的注册表能判断任务是否在跑。 进程重启后注册表为空,最终要以数据库状态为准。
- 误区:进度到 100 就说明完成了。 进度和状态是两个字段,完成要看状态;进度中途封顶,收尾时才置 100。
- 误区:SSE 连接等浏览器自己断开就行。 关弹窗、离开页面、发出中断时都要主动断开,否则连接一直占着。
- 追问:轮询间隔怎么定? 看阶段变化的频率:一次模型调用要好几秒,2–3 秒一次足够;没有运行中的任务时就不要发。
- 追问:SSE 推增量还是推全量? 任务数据不大时推全量,前端整体替换最不容易错;数据很大时才考虑增量。
九、加强记忆
耗时流程异步化记「交、跑、看」:交是接口建记录、交给后台、立即返回任务 ID;跑是后台任务四件事,注册表防重复调度且结束即摘除、信号量限流让其余任务排队、最外层兜住异常把终态落库、用户和类加载器这些上下文显式带进去并在 finally 清理。看是进度等于阶段描述加单调递增的百分比,按区间或步数计算、中途封顶、收尾置 100;列表页轮询、只在有运行中的行时发请求,详情页用 SSE 每秒读库、有变化才推、结束状态断开,前端按两个换行拼报文并在离开时主动断开。
项目实战落地
项目里怎么做的
《AI Agent智能会议纪要辅助系统》的转写、纪要生成、纪要自检三类任务全部在后台异步执行:
- 调度:以自检为例,
schedule_agent_run先查_running_agents,同一次运行已经在跑就直接返回;否则asyncio.create_task创建任务放进字典,add_done_callback在结束时把自己摘掉;发起自检、恢复运行、服务启动续跑都走这个函数; - 限流:转写和纪要生成各用一个
Semaphore(2),第一次用到时才创建;转写的第三个任务停在async with处排队,状态仍是 PENDING,页面显示「等待执行」; - 进度:每开一个步骤就更新
stage和progress;自检的进度是min(95, 步号 × 100 / (max_rounds × 2 + 1)),最后 5 个点留给汇总步结束时置 100; - SSE:详情弹窗打开运行中的记录时连
GET /agent/events/{run_id},服务端每秒读一次运行记录,内容和上次推的不同才推,进入成功、失败、中断、放弃、等待确认五种状态后结束循环;前端用fetch带token请求头读流,按\n\n切报文,半条放回缓冲区,拿到后整体替换弹窗数据。
《AI Agent岗位匹配与求职规划系统》用的是 Java 线程池加前端轮询:CompletableFuture.runAsync 把 ReAct Agent 交给 ForkJoinPool 公共线程池,HTTP 请求先返回;后台线程第一件事是把上下文类加载器换成本项目的,否则打成 jar 运行时,向量模型类初始化失败抛出的是 Error,catch (Exception e) 接不住,运行记录会一直停在运行中。
为什么这样取舍
- SSE 替代轮询看详情:省掉前端定时轮询,后台每写一次运行记录或步骤,下一秒的查询结果就和上次不同,接口推出一条新报文。
- 进度封顶 95:一轮两步加最后一步汇总是整次运行最多的步数,封顶 95 免得进度条提前走满。
面试官还会追问
- 服务器上没装 FFmpeg 时,转写任务会怎样?页面上能看到什么?
- SSE 接口在建立事件流之前先做了什么校验?员工能连上别人发起的自检事件流吗?
- 一个模型配置被自检运行用过之后,管理员还能改它的用途或删掉它吗?为什么?
学完《AI Agent智能会议纪要辅助系统》,上面这些追问你都会迎刃而解。