Python 异步任务队列:Celery + Redis/RabbitMQ 实战

Celery 分布式任务队列实战:任务定义、定时调度(Beat)、结果存储、监控(Flower)、重试策略、死信队列与生产最佳实践。

1. Celery 架构

Application
    ├── Producer(Flask/Django/FastAPI)──→ 发送任务
    │
    ├── Broker(Redis / RabbitMQ)───────→ 任务队列存储
    │
    ├── Worker(多进程/多线程)────────→ 消费并执行任务
    │
    └── Backend(Redis / Database)────→ 存储任务结果

2. 快速上手

from celery import Celery

# 创建 Celery 应用
app = Celery('tasks',
    broker='redis://localhost:6379/0',
    backend='redis://localhost:6379/1'
)

# 定义任务
@app.task(bind=True, max_retries=3)
def send_email(self, to: str, subject: str, body: str):
    """发送邮件任务"""
    try:
        smtp.send(to, subject, body)
        return {"status": "sent", "to": to}
    except Exception as exc:
        # 指数退避重试
        raise self.retry(exc=exc, countdown=2 ** self.request.retries)

@app.task
def process_image(image_path: str):
    """处理图片任务"""
    from PIL import Image
    img = Image.open(image_path)
    img.thumbnail((800, 600))
    img.save(image_path + '.thumb.jpg')
    return {"thumbnail": image_path + '.thumb.jpg'}

# 定时任务
@app.on_after_configure.connect
def setup_periodic_tasks(sender, **kwargs):
    sender.add_periodic_task(60.0, cleanup.s(), name='cleanup every minute')
    sender.add_periodic_task(crontab(hour=2, minute=0), backup.s())

3. 调用任务

# 异步调用(立即返回.AsyncResult)
result = send_email.delay('user@example.com', 'Hello', 'Welcome!')
print(result.id)  # 获取任务 ID

# 查询结果
result = send_email.AsyncResult(result.id)
if result.ready():
    print(result.get(timeout=10))  # 获取结果

# 定时执行
send_email.apply_async(
    args=['user@example.com', 'Hello', 'Welcome!'],
    countdown=300,      # 5 分钟后执行
    eta=datetime(2024, 1, 1, 10, 0),  # 指定时间
    queue='email'       # 指定队列
)

4. 生产配置

# celeryconfig.py
broker_url = 'amqp://user:password@rabbitmq:5672/myvhost'
result_backend = 'redis://redis:6379/0'

# 序列化
task_serializer = 'json'
result_serializer = 'json'
accept_content = ['json']

# Worker 配置
worker_prefetch_multiplier = 1  # 一次只取一个任务
worker_max_tasks_per_child = 1000  # 执行 N 个任务后重启 Worker

# 任务结果过期时间
result_expires = 3600  # 1 小时

# 路由
task_routes = {
    'tasks.send_email': {'queue': 'email'},
    'tasks.process_image': {'queue': 'image'},
}

# 队列定义
task_queues = (
    Queue('default', routing_key='default'),
    Queue('email', routing_key='email'),
    Queue('image', routing_key='image'),
)

5. 监控与运维

# 启动 Worker
celery -A tasks worker --loglevel=info --queues=default,email,image -c 4

# 启动 Beat(定时任务调度)
celery -A tasks beat --loglevel=info

# Flower 监控界面
celery -A tasks flower --port=5555
# 访问 http://localhost:5555 查看任务状态、Worker 负载、队列深度

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「python」更多文章

  1. Python 高级异步编程:Trio 结构化并发与 AnyIO 兼容层
  2. Python 数据工程与 ETL 管道实战
  3. Python 元编程与动态特性深度解析