← 工作流编排

耗时的 AI 流程怎么异步执行?进度怎么推给前端?

高频 中等 长任务异步执行与进度推送 · 第 1 / 1 问 更新于 2026/09/29
工作流异步任务SSE轮询并发控制
本题落地项目AI Agent智能会议纪要辅助系统

简化版

接口只负责建一条任务记录(状态「运行中」或「待执行」)并把任务交给后台,立刻返回任务 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智能会议纪要辅助系统》,上面这些追问你都会迎刃而解。

本题落地项目地狱锤炼AI Agent智能会议纪要辅助系统基于LangChain+LangGraph的AI Agent智能会议纪要辅助系统 包括会议管理 资料管理 音视频转写 说话人分离 会议纪要生成 Agent自检更正 Agent运行观测等。FastAPILangChainLangGraphAgent多模态源码+SQL喂饭学习教程配套面试文档环境安装文档项目运行文档 学习这个项目 也可以学AI Agent岗位匹配与求职规划系统地狱锤炼 查看项目