共计 2200 个字符,预计需要花费 6 分钟才能阅读完成。
为什么需要 Agent 系统?
刚开始尝试写 Agent 程序时,我直接用 Python 的 threading 模块裸写多线程任务,结果踩了不少坑:

- 任务执行到一半进程崩溃,状态全丢失
- 多个 Worker 同时抢同一个任务导致重复处理
- 想查看任务历史记录时,只能翻日志文件
最痛苦的是有次处理支付回调,因为线程阻塞导致漏单,不得不手工补数据。这才意识到需要更专业的任务调度方案。
技术选型:Celery vs RQ vs Airflow
调研了主流方案后,对比结果如下:
- Celery:分布式任务队列标杆,支持 Redis/RabbitMQ 等多种 broker,适合需要弹性扩展的场景
- RQ (Redis Queue):基于 Redis 的轻量级方案,学习曲线低但功能较单一
- Airflow:面向复杂工作流设计,调度功能强大但太重
最终选择 Celery 的原因:
- 我们的支付回调场景需要秒级任务重试
- 未来可能需要水平扩展 Worker
- 自带 flower 监控面板方便运维
基础架构搭建
安装与最小化配置
先安装核心组件:
pip install celery redis flower
创建 tasks.py 定义基础结构:
from celery import Celery
app = Celery(
'agent_demo',
broker='redis://localhost:6379/0',
backend='redis://localhost:6379/1'
)
@app.task(bind=True)
def process_payment(self, order_id: str) -> bool:
"""处理支付回调任务"""
try:
# 业务逻辑代码
return True
except Exception as e:
self.retry(exc=e, countdown=60) # 60 秒后重试
关键配置解析
在 celeryconfig.py 中添加生产环境配置:
task_acks_late = True # 任务完成后才确认
result_expires = 3600 # 结果保存 1 小时
worker_prefetch_multiplier = 1 # 防止 Worker 囤积任务
可靠性增强实战
智能重试机制
通过装饰器实现分级重试策略:
from celery.exceptions import Retry
@app.task(
bind=True,
autoretry_for=(ConnectionError, TimeoutError),
retry_backoff=True,
max_retries=3
)
def fetch_api_data(self, url: str):
"""
带指数退避的 API 请求任务
:param url: 目标 API 地址
:raises: 非网络错误直接失败
"""
response = requests.get(url, timeout=10)
response.raise_for_status() # 非 2xx 状态码触发重试
return response.json()
状态监控方案
集成 Flask-Admin 的完整示例:
from flask import Flask
from flask_admin import Admin
from celery.contrib import admin as celery_admin
flask_app = Flask(__name__)
admin = Admin(flask_app, name='Agent 监控')
admin.add_view(celery_admin.TaskModelView(app)) # 注入 Celery 监控
访问 /admin 即可看到实时任务状态:
- 成功 / 失败计数
- 当前执行中的 Worker
- 历史任务参数和结果
生产环境避坑指南
性能调优
Worker 配置公式:
推荐 Worker 数 = CPU 核心数 × 2 + 1
但要注意:
- I/ O 密集型任务可适当增加
- 每个 Worker 默认会预取 4 倍任务(通过
prefetch_multiplier调整)
序列化陷阱
常见踩坑点:
- 传递了数据库连接对象等不可 pickle 的参数
- 自定义异常未在 Worker 端注册导致反序列化失败
解决方案:
# 自定义异常必须定义在共享模块中
class PaymentError(Exception):
pass
# 复杂对象建议传递 ID 而非实例
@app.task
def sync_order(order_id: str): # ✅ 好习惯
order = Order.get(order_id)
...
进阶思考:任务依赖链
当需要实现 ” 任务 A 成功后再执行任务 B ” 时,可以:
- 在 A 的
on_success回调中触发 B(紧耦合) - 使用 Celery 的
chain组合任务 - 复杂场景建议用
chord实现并行任务 + 回调
示例代码:
from celery import chord
header = [fetch_data.s(url) for url in url_list]
callback = aggregate_results.s()
chord(header)(callback) # 所有 header 任务完成后执行 callback
总结与资源推荐
经过这次实践,我的 Agent 系统终于可以:
- 自动恢复中断的任务
- 实时监控执行状态
- 平滑应对流量高峰
推荐继续学习:
- 官方文档《Monitoring and Management Guide》
- 源码阅读
celery.app.task.Task类 - 进阶书《Celery in Action》
下次可以尝试用 Kubernetes 部署 Celery 集群,实现真正的弹性伸缩。
正文完
