Python 消息队列与任务系统:RabbitMQ、Redis、Celery

集中整理 Python 项目里的消息与任务系统主线:RabbitMQ 负责可靠投递,Redis 负责缓存与轻量结构,Celery 负责组织异步任务流程。

这篇文章只讲一件事:当程序开始有异步任务、消息投递和任务调度时,各个组件应该怎么分工。

一、先区分 4 个常见角色

  • RabbitMQ:专业消息队列,适合异步任务和系统解耦
  • Redis:内存数据库,常用于缓存、计数和轻量结构
  • Celery:分布式任务队列框架,不是消息队列本身
  • APScheduler:轻量调度器,适合单机或简单定时任务

可以先记住:Celery 依赖消息中间件或结果后端工作,它自己不是底层消息系统。

二、为什么消息队列不是普通 queue.Queue

Python 里的 queue.Queue 和 RabbitMQ 不是一个层级的东西。

本地队列 queue.Queue

  • 只能在当前 Python 进程体系里用
  • 更适合线程间解耦
  • 不能天然跨机器、跨语言

RabbitMQ

  • 是独立的消息中间件服务
  • 可以跨机器、跨语言传递消息
  • 更适合系统级异步解耦

三、RabbitMQ 的最小生产消费模型

生产者

import pika

connection = pika.BlockingConnection(
    pika.ConnectionParameters("localhost")
)
channel = connection.channel()

channel.queue_declare(queue="hello")
channel.basic_publish(
    exchange="",
    routing_key="hello",
    body="Hello World!",
)

print("message sent")
connection.close()

消费者

import pika

connection = pika.BlockingConnection(
    pika.ConnectionParameters("localhost")
)
channel = connection.channel()
channel.queue_declare(queue="hello")


def callback(ch, method, properties, body):
    print("received:", body.decode("utf-8"))
    ch.basic_ack(delivery_tag=method.delivery_tag)


channel.basic_consume(queue="hello", on_message_callback=callback)
print("waiting for messages...")
channel.start_consuming()

这里真正重要的知识点

  • 消息不是直接“发给队列”,而是经过 exchange 路由
  • 这个例子里的 exchange="" 表示使用 RabbitMQ 的默认交换机,它会按 routing_key 把消息路由到同名队列
  • 生产者和消费者都可以重复声明队列
  • 手动 ack 能降低消费者异常退出导致的丢消息风险

四、RabbitMQ 的几个高频概念

1. 持久化

队列持久化:

channel.queue_declare(queue="task_queue", durable=True)

消息持久化:

channel.basic_publish(
    exchange="",
    routing_key="task_queue",
    body="message",
    properties=pika.BasicProperties(delivery_mode=2),
)

2. 公平分发

channel.basic_qos(prefetch_count=1)

3. Exchange

常见分发模式可以概括为:

  • 直连
  • 广播
  • 路由
  • 主题匹配

五、Redis 在这条主线里的定位

Redis 更适合:

  • 缓存热点数据
  • 保存 Session 或验证码状态
  • 做计数器和排行榜
  • 提供轻量消息结构或延迟任务辅助结构

位图统计在线人数

旧资料里提到的 setbit / getbit / bitcount 很有代表性。

例如:

  • 用户上线:SETBIT online_users 200 1
  • 判断用户是否在线:GETBIT online_users 200
  • 统计当前在线总人数:BITCOUNT online_users

适合场景:

  • 用户 ID 连续或可映射为偏移量
  • 只需要布尔状态统计

六、Redis 和 RabbitMQ 怎么选

可以先这样判断:

  • 要专业消息投递、确认机制、路由能力:优先 RabbitMQ
  • 要高速缓存、计数器、排行榜、轻量结构:优先 Redis

不要因为两者都能“放数据”,就把它们当成同一类组件。

七、Celery 的真正位置

Celery 常用于:

  • 异步发送邮件
  • 异步生成报表
  • 延迟执行任务
  • 周期性任务调度

Celery 依赖什么

Celery 通常需要:

  • broker:任务消息中间件,例如 Redis 或 RabbitMQ
  • worker:真正执行任务的进程
  • result backend:保存任务结果的地方,可选

一个最小任务示例

from celery import Celery

app = Celery(
    "demo",
    broker="redis://127.0.0.1:6379/0",
    backend="redis://127.0.0.1:6379/1",
)


@app.task
def add(x: int, y: int) -> int:
    return x + y

调用时:

result = add.delay(1, 2)
print(result.id)

还需要启动 worker,任务才会真正被消费:

celery -A tasks worker -l info

如果你的文件名不是 tasks.py,把这里的 tasks 改成实际模块名即可。

八、定时任务不一定非得上 Celery

如果你的需求只是:

  • 每隔几分钟跑一次脚本
  • 每天固定时间做同步
  • 在单机服务里维护少量计划任务

那并不一定要直接上 Celery + BeatAPScheduler 往往更轻量。

from apscheduler.schedulers.blocking import BlockingScheduler


def sync_report() -> None:
    print("run scheduled job")


scheduler = BlockingScheduler()
scheduler.add_job(sync_report, "cron", hour=2, minute=0)
scheduler.start()

可以先这样判断:

  • 重点是“定时调度”且规模较小:先看 APScheduler
  • 重点是“分布式异步任务 + worker 消费”:再看 Celery

九、轻量替代方案

如果项目不需要 Celery 的完整功能,也可以关注:

框架 特点 适用场景
dramatiq API 简洁,性能好 中小型项目
rq 基于 Redis,极简 简单异步任务
arq 基于 asyncio 异步项目
huey 轻量,支持 Redis/SQLite 小型项目

十、实战建议

  • 先分清“缓存”和“可靠消息”不是同一件事
  • 先学清楚 RabbitMQ 的确认机制,再谈消费可靠性
  • 用 Redis 时,不要把它简单理解成“什么都能塞”
  • 先确认需求是“调度”还是“分布式异步任务”,再决定是否上 Celery

小结

把这条主线压缩成一句话就是:

  • RabbitMQ 负责可靠消息流转
  • Redis 负责高速访问与轻量结构
  • Celery 负责把异步任务流程组织起来

理解这套分工后,再看后台系统、运维平台和 Web 服务里的任务链路,会更有全局感。

相关阅读

  • 《Python 网络与并发》:理解线程、进程、本地队列与高并发 IO
  • 《Python 数据访问:PyMySQL 与 SQLAlchemy 入门》:数据库访问是另一条并行主线