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),前台照样秒速接单。今天全程用这套"后厨"类比。
这就像生活中的餐厅后厨:前台(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 的前台——服务员写下
① 某个空闲的厨师(worker)撕下你,照着单号去 Postgres 查这桌客人做到哪了(current_step_id);
② 他只做一道工序:搭提示词→问 LLM→按回答动一次手(调工具),把过程记在 Feed 上;
③ 做完他看一眼状态:还没上齐?他就手写一张新小票(apply_async, 2 秒后生效)插回夹子——那是你的"下一步兄弟";
④ 而你,此刻已功成身退。整桌菜就是被一张张短命小票接力做完的。
(execution_id=88, 时间) 就把你插进了 Redis 小票夹。① 某个空闲的厨师(worker)撕下你,照着单号去 Postgres 查这桌客人做到哪了(current_step_id);
② 他只做一道工序:搭提示词→问 LLM→按回答动一次手(调工具),把过程记在 Feed 上;
③ 做完他看一眼状态:还没上齐?他就手写一张新小票(apply_async, 2 秒后生效)插回夹子——那是你的"下一步兄弟";
④ 而你,此刻已功成身退。整桌菜就是被一张张短命小票接力做完的。
L03
自我重排循环(核心中的核心)
AgentExecutor.execute_next_step(jobs/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)
没有大 while 循环——每跑一步就给自己排下一步,进度全落 DB。所以进程可断可恢复、可多开、崩溃不丢进度。
📝 一次 3 步的运行,队列里发生了什么
全程没有任何一个进程"一直占着"——每步都是独立的短任务,中间断电重启也能从 DB 里的 step_id 接着跑。
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 决定——status 不是 RUNNING 就不再排小票,循环自然停。
| 时刻 | 谁在干什么 | 队列里 | DB: status / step |
|---|---|---|---|
| T=0s | API 调 delay() | 1 张小票 | RUNNING / step1 |
| T=0.1s | worker 领走,跑第 1 步(LLM+工具) | 空 | RUNNING / step1 |
| T=30s | 第 1 步完,apply_async(countdown=2) | 1 张(2 秒后可领) | RUNNING / step2 |
| T=32s | worker 领走,跑第 2 步 | 空 | RUNNING / step2 |
| T=61s | 第 2 步完,再排一张 | 1 张 | RUNNING / step3 |
| T=90s | 第 3 步跑完,目标达成 → 不再排 | 空(永远) | COMPLETED / — |
"循环 = 任务反复自我重新调度"
这是 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=5;execute_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-agent 调 AgentScheduleHelper.run_scheduled_agents(helper/agent_schedule_helper.py:16)——查出"过去 5 分钟内到点、status=SCHEDULED"的计划,交给 ScheduledAgentExecutor(jobs/scheduling_executor.py:25):新建一条 AgentExecution、设 RUNNING、execute_agent.delay(...)。
定时和手动殊途同归
定时触发和用户手动点"运行"最后都调
execute_agent.delay()——汇入同一个执行任务。Scheduler 只是"按时间自动调用"那个入口,执行逻辑完全复用(和 AutoGPT Day 14 一样)。这让"每天早上自动跑智能体"成为可能——SuperAGI"持续运行的智能体"(Day 01)的价值就在于此。"多种触发源 → 同一执行入口"是清晰架构的标志。L07
其它后台任务
summarize_resource(worker.py:75):上传文件后异步做摘要 + 写向量库(Day 14)。webhook_callback(worker.py:106):智能体状态变化时通过 SQLAlchemyevent.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 18:GUI 与前端——Next.js 界面怎么创建智能体、实时监控执行、逛工具市场。前端如何驱动整个后端。