Day 17 / 共 20 天 · 第 4 周 平台/异步/生态

Celery 异步执行

平台的核心机制。智能体的自主循环在 Celery 后台 worker 里跑(不阻塞 API)。今天读 worker.py + jobs/,看"循环在队列里"如何实现。

📍 你在整门课的位置 · 第 4 周 平台/异步/生态
D16 API(把运行入队) D17 Celery(队列里跑循环) D18 GUI 看进度 D19 部署 D20 收官
L01

为什么异步

🤔 Day 16 说"API 把任务 delay 出去就秒回"——那到底是谁在跑智能体? 就是今天的主角:Celery worker。API 只是把一张"运行单"塞进 Redis 队列;真正埋头干活(一轮轮调 LLM、执行工具)的是后台的 worker 进程。今天就看这张单子怎么被领走、怎么一步步把循环跑完。

一个自主智能体要循环很多轮(每轮调 LLM、可能几十秒),绝不能卡在 HTTP 请求里。所以 SuperAGI 把"运行智能体"变成 Celery 后台任务。

💡 一句话本质 Celery = "生产者 / 队列 / 消费者" 三件套:API(生产者)把任务丢进 Redis(队列),worker(消费者)领走执行。三者解耦——所以能秒回、能多开 worker 扛并发、能崩溃后重来。
这就像生活中的餐厅后厨:前台(API)接单开小票、把小票插进订单小票夹(Redis 队列)、厨师(worker)从夹子上撕单做菜。前台不做菜、厨师不接单、小票夹两头解耦——生意忙了就多雇几个厨师(多开 worker),前台照样秒速接单。今天全程用这套"后厨"类比。
同步 vs 异步的天壤之别 如果智能体在 HTTP 请求里同步跑:用户点"运行",浏览器就转圈等几分钟(甚至超时断开)、服务器一个请求占用几分钟。异步:API 收到"运行"请求,把任务丢进 Celery 队列、秒回;后台 worker 慢慢跑;用户轮询看进度。API 保持轻快、能扛高并发,长任务在后台可靠执行。对"长时间自主任务",异步是唯一正确选择——AutoGPT(MQ+executor)、OpenHands(双进程)都这么做。
L02

execute_agent 任务

worker.py:66——智能体的每一步是一个 Celery task:

@app.task(name="execute_agent", autoretry_for=(Exception,), retry_backoff=2, max_retries=5)
def execute_agent(agent_execution_id: int, time):
    from superagi.jobs.agent_executor import AgentExecutor
    handle_tools_import()
    AgentExecutor().execute_next_step(agent_execution_id=agent_execution_id)
读法:Day 16 的 execute_agent.delay(...) 就是把这个任务推给 Redis 队列,某个 worker 领走执行。注意它叫 execute_next_step下一步)——不是"跑到底",只跑一步(Day 06)。broker 和 backend 都指向 Redis(worker.py:25)。
⚠️ 小白常误以为:worker 领到任务后会开一个 while True 把智能体从头跑到底。其实:一个 execute_agent 任务只跑一步就结束——"跑很多步"是靠 L03 的自我重排把很多个一步串起来的。
📝 第一人称之旅:现在你是一张 execute_agent 小票 你出生在 Day 16 的前台——服务员写下 (execution_id=88, 时间) 就把你插进了 Redis 小票夹。
① 某个空闲的厨师(worker)撕下你,照着单号去 Postgres 查这桌客人做到哪了(current_step_id);
② 他只做一道工序:搭提示词→问 LLM→按回答动一次手(调工具),把过程记在 Feed 上;
③ 做完他看一眼状态:还没上齐?他就手写一张新小票(apply_async, 2 秒后生效)插回夹子——那是你的"下一步兄弟";
④ 而你,此刻已功成身退。整桌菜就是被一张张短命小票接力做完的。
L03

自我重排循环(核心中的核心)

AgentExecutor.execute_next_stepjobs/agent_executor.py:39)跑完一步,若没完成就再排一个任务给自己:94):

worker 跑第 N 步 apply_async(countdown=2) Redis 队列 worker 跑第 N+1 步
if agent_execution.status in ("COMPLETED", "WAITING_FOR_PERMISSION"):
    return
superagi.worker.execute_agent.apply_async((agent_execution_id, datetime.now()), countdown=2)
worker 跑第 N 步 apply_asynccountdown=2 Redis 队列排下一步任务 worker 领走第 N+1 步 状态存 DBstep_id + Feed status=COMPLETED / WAITING_FOR_PERMISSION 时不再重排 → 循环停
没有大 while 循环——每跑一步就给自己排下一步,进度全落 DB。所以进程可断可恢复、可多开、崩溃不丢进度。
📝 一次 3 步的运行,队列里发生了什么 T=0s worker 领任务→跑第1步(想+调工具)→写 Feed、存 step_id→status 仍 RUNNING→apply_async(countdown=2)
T=2s 领第2步…同样再排一个→T=4s 领第3步→这步 status 变 COMPLETED不再重排,循环自然结束
全程没有任何一个进程"一直占着"——每步都是独立的短任务,中间断电重启也能从 DB 里的 step_id 接着跑。
📝 单步走查表:一次 3 步运行的每个瞬间
时刻谁在干什么队列里DB: status / step
T=0sAPI 调 delay()1 张小票RUNNING / step1
T=0.1sworker 领走,跑第 1 步(LLM+工具)RUNNING / step1
T=30s第 1 步完,apply_async(countdown=2)1 张(2 秒后可领)RUNNING / step2
T=32sworker 领走,跑第 2 步RUNNING / step2
T=61s第 2 步完,再排一张1 张RUNNING / step3
T=90s第 3 步跑完,目标达成 → 不再排空(永远)COMPLETED / —
盯住最后一列:"循环还继不继续"完全由 DB 里的 status 决定——status 不是 RUNNING 就不再排小票,循环自然停。
"循环 = 任务反复自我重新调度" 这是 SuperAGI 最精妙的设计(Day 02/05 反复强调):不是一个进程里的大 while 循环,而是"每跑一步就排下一步任务"——Celery 任务链自我延续。就像后厨做一桌宴席:厨师不会守着一桌菜从头做到尾,而是每做完一道,就写一张"下一道菜"的小票插回夹子——哪个厨师有空谁接着做,做到哪一道全记在单子上(DB),换厨师、下班交接都不乱。进度存 DB(current_agent_step_id)、对话存 Feed。好处:① 进程可断可恢复(状态都持久化);② 多 worker 能水平扩展;③ 崩溃不丢进度;④ 天然支持暂停/审批挂起(status 变了就不再重排)。countdown=2 让每步之间隔 2 秒——避免打爆 LLM API。
L04

自动重试

execute_agent 任务本身配了 autoretry_for=(Exception,), retry_backoff=2, max_retries=5execute_next_step 出错时也 apply_async(countdown=15) 延迟重试(agent_executor.py:89)。

两层重试保障 智能体一步可能因各种原因失败(LLM 限流、工具报错、网络抖动)。Celery 的 autoretry_for 自动重试(最多 5 次、指数退避),加上出错时延迟 15 秒再排——让智能体面对临时故障能自愈,不会一遇错就死。加上 Day 10 的 LLM 调用层重试,多层容错保证长任务的韧性。"长时间运行的东西必须能扛住临时故障"——异步任务队列 + 自动重试是标准答案。
L05

Celery beat 定时

worker.py:32 配置了 Celery beat 定时任务(Day 04 的 --beat):

beat_schedule = {
  'initialize-schedule-agent': {'schedule': timedelta(minutes=5)},   # 每 5 分钟
  'execute_waiting_workflows': {'schedule': timedelta(minutes=2)},   # 每 2 分钟
}
读法:Celery beat 是"定时闹钟"——按固定间隔触发任务。initialize-schedule-agent(每 5 分钟)检查有没有到点的定时智能体(L06);execute_waiting_workflows(每 2 分钟)唤醒等待中的工作流(WAIT_STEP 时间到了的)改回 RUNNING 并重排。
还是那间后厨 Celery beat 就像后厨领班的两个定时闹钟:每 5 分钟翻一遍预约单("有没有客人订了 12 点的席?到点开做"= 定时智能体);每 2 分钟看一眼醒发架("哪份面团醒够时间了?拿回灶台"= 唤醒 WAIT_STEP 到时的工作流)。闹钟只负责"到点提醒",做菜仍然走同一个小票流程。
L06

定时智能体

initialize-schedule-agentAgentScheduleHelper.run_scheduled_agentshelper/agent_schedule_helper.py:16)——查出"过去 5 分钟内到点、status=SCHEDULED"的计划,交给 ScheduledAgentExecutorjobs/scheduling_executor.py:25):新建一条 AgentExecution、设 RUNNING、execute_agent.delay(...)

定时和手动殊途同归 定时触发和用户手动点"运行"最后都调 execute_agent.delay()——汇入同一个执行任务。Scheduler 只是"按时间自动调用"那个入口,执行逻辑完全复用(和 AutoGPT Day 14 一样)。这让"每天早上自动跑智能体"成为可能——SuperAGI"持续运行的智能体"(Day 01)的价值就在于此。"多种触发源 → 同一执行入口"是清晰架构的标志。
L07

其它后台任务

  • summarize_resourceworker.py:75):上传文件后异步做摘要 + 写向量库(Day 14)。
  • webhook_callbackworker.py:106):智能体状态变化时通过 SQLAlchemy event.listens_for 触发 webhook 回调(通知外部系统)。
为什么这些也放后台? 文件摘要要调 LLM(慢)、webhook 要发 HTTP(可能慢/失败)——都不该阻塞主流程。凡是"慢的、可失败的、不需即时结果的"操作,都丢给 Celery 后台异步跑。主流程只管快速响应,副作用后台处理。Celery 是平台所有异步工作的统一执行器——智能体循环、定时、摘要、webhook 全走它。
L08

今日小结 + 动手

🧠 今天你应该能回答

  • 为什么智能体执行必须异步?同步会怎样?
  • execute_agent 任务为什么叫"next_step"?
  • "自我重排循环"怎么工作?好处(可断可恢复/可扩展/不丢进度)?
  • 两层重试保障是什么?
  • Celery beat 定时干什么?定时和手动为什么殊途同归?

✋ 动手

P=superagi
sed -n '25,72p' $P/worker.py | head -40             # Celery app + beat + execute_agent
sed -n '39,100p' $P/jobs/agent_executor.py | head -40   # 自我重排
grep -n 'def run_scheduled_agents' $P/helper/agent_schedule_helper.py
sed -n '25,83p' $P/jobs/scheduling_executor.py | head -30
明天预告 · Day 18GUI 与前端——Next.js 界面怎么创建智能体、实时监控执行、逛工具市场。前端如何驱动整个后端。
← Day 16 API Day 18 · GUI →