异步工作流与后台任务
长期运行的智能体需要队列(Celery、RQ、Temporal)——返回任务 ID,并轮询状态。
异步工作流与后台任务 是 CoddyKit 上的免费 AI Agents 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 AI Agents 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 AI Agents 课程共包含 4 节课。
本课时的部分内容尚未翻译,以英文显示。
请求周期过慢时
某些代理任务需要 30 秒、5 分钟甚至数小时。您不能让 HTTP 连接一直保持打开状态。请将它们移至异步后台作业中。
模式
- 客户端通过 POST 提交任务 -> 服务器创建作业并返回 job_id
- 工作进程从队列中获取作业
- 客户端轮询 GET /jobs/{id} 获取状态
- 完成后,服务器返回结果
Submit Endpoint
from uuid import uuid4
@app.post('/jobs')
def submit(req: AgentRequest):
job_id = str(uuid4())
redis.hset(f'job:{job_id}', mapping={'status': 'queued', 'user_id': req.user_id})
queue.enqueue('run_agent_job', job_id, req.query)
return {'job_id': job_id, 'status': 'queued'}Status Endpoint
@app.get('/jobs/{job_id}')
def status(job_id: str):
data = redis.hgetall(f'job:{job_id}')
if not data:
raise HTTPException(404)
return {
'status': data['status'],
'result': data.get('result'),
'error': data.get('error')
}工作进程
一个独立进程从队列中消费作业:
def run_agent_job(job_id, query):
redis.hset(f'job:{job_id}', 'status', 'running')
try:
result = run_agent(query)
redis.hset(f'job:{job_id}', mapping={'status': 'done', 'result': result})
except Exception as e:
redis.hset(f'job:{job_id}', mapping={'status': 'failed', 'error': str(e)})队列选择
- RQ(Redis Queue)——极简的 Python 队列
- Celery——经典的 Python 方案,功能丰富
- Temporal——持久化工作流、重试和可观测性
- Dramatiq——Celery 的替代方案
- 云原生——Cloud Tasks、SQS、Pub/Sub
使用 Temporal 实现持久化工作流
代理工作流是有状态的。Temporal 的持久化执行模型非常适合这种场景:
import temporalio
@temporalio.workflow.defn
class AgentWorkflow:
@temporalio.workflow.run
async def run(self, query: str) -> str:
plan = await workflow.execute_activity(plan_step, query, schedule_to_close_timeout=timedelta(minutes=2))
results = await workflow.execute_activity(execute_step, plan, schedule_to_close_timeout=timedelta(minutes=10))
return await workflow.execute_activity(synthesise_step, results)流式传输部分更新
长时间运行的作业可以从进度更新中受益。请使用服务器发送事件或 WebSockets:
@app.get('/jobs/{job_id}/stream')
def stream_progress(job_id):
def gen():
while True:
update = redis.brpop(f'updates:{job_id}', timeout=30)
if not update:
yield 'data: {"status": "timeout"}\n\n'
break
yield f'data: {update[1].decode()}\n\n'
if 'done' in update[1].decode():
break
return StreamingResponse(gen(), media_type='text/event-stream')作业存活时间
不要永久保留作业记录:
redis.expire(f'job:{job_id}', 86400) # 1 day重试
临时故障会自动重试;永久性故障会进入 DLQ,供人工审核:
import time
def retry(retries=3, retry_backoff=True):
def decorator(func):
def wrapper(*args, **kwargs):
for attempt in range(1, retries + 1):
try:
return func(*args, **kwargs)
except Exception as e:
print(f'attempt {attempt} failed: {e}')
if attempt == retries:
raise
return None
return wrapper
return decorator
attempts = {'n': 0}
@retry(retries=3, retry_backoff=True)
def run_agent_job(job_id):
attempts['n'] += 1
if attempts['n'] < 3:
raise RuntimeError('transient error')
return f'job {job_id} done'
print(run_agent_job('job-1'))
并发限制
为了控制 GPU 和接口成本,请限制并发工作进程数:
rq worker --burst --max-jobs 1000 --queue agent
# Or per-queue concurrency in Temporal worker config.每位用户的配额
请跟踪每位用户正在处理的作业:
key = f'inflight:{user_id}'
if redis.scard(key) >= 5:
raise HTTPException(429, 'Too many in-flight jobs')
redis.sadd(key, job_id)可观测性
端到端跟踪每个作业。使用 job_id 和 user_id 标记各个跨度。失败的作业应自动创建工单或触发告警。
状态轮询模式
为什么长时间运行的代理要使用提交后轮询状态的模式?
回顾
通过 POST 提交作业,返回 id,由工作进程从队列中处理,客户端进行轮询。简单场景使用 RQ/Celery;有状态工作流使用 Temporal。还需要设置存活时间、重试、配额和跟踪。
常见问题解答
「异步工作流与后台任务」课时是免费的吗?
是的 — 「异步工作流与后台任务」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 AI Agents 课程的其余内容,请升级到 CoddyKit PRO。 AI Agents 课程共包含 4 节课。
「异步工作流与后台任务」这节课中我会学到什么?
长期运行的智能体需要队列(Celery、RQ、Temporal)——返回任务 ID,并轮询状态。 你通过在浏览器中直接运行的动手代码来练习 AI Agents,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 AI Agents 需要有经验吗?
无需任何先前经验。CoddyKit 上的 AI Agents 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。
「异步工作流与后台任务」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 AI Agents 课中编写并运行代码吗?
能。每节 AI Agents 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- 通过 API 提供智能体服务
- 异步工作流与后台任务
- 速率限制与配额管理
- 智能体的蓝绿部署与金丝雀部署