Agent小项目实战:从零搭建高可用的自动化任务系统

1次阅读
没有评论

共计 2200 个字符,预计需要花费 6 分钟才能阅读完成。

image.webp

为什么需要 Agent 系统?

刚开始尝试写 Agent 程序时,我直接用 Python 的 threading 模块裸写多线程任务,结果踩了不少坑:

Agent 小项目实战:从零搭建高可用的自动化任务系统

  • 任务执行到一半进程崩溃,状态全丢失
  • 多个 Worker 同时抢同一个任务导致重复处理
  • 想查看任务历史记录时,只能翻日志文件

最痛苦的是有次处理支付回调,因为线程阻塞导致漏单,不得不手工补数据。这才意识到需要更专业的任务调度方案。

技术选型:Celery vs RQ vs Airflow

调研了主流方案后,对比结果如下:

  • Celery:分布式任务队列标杆,支持 Redis/RabbitMQ 等多种 broker,适合需要弹性扩展的场景
  • RQ (Redis Queue):基于 Redis 的轻量级方案,学习曲线低但功能较单一
  • Airflow:面向复杂工作流设计,调度功能强大但太重

最终选择 Celery 的原因:

  1. 我们的支付回调场景需要秒级任务重试
  2. 未来可能需要水平扩展 Worker
  3. 自带 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 ” 时,可以:

  1. 在 A 的 on_success 回调中触发 B(紧耦合)
  2. 使用 Celery 的 chain 组合任务
  3. 复杂场景建议用 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 集群,实现真正的弹性伸缩。

正文完
 0
评论(没有评论)