这篇文章只讲一件事:当程序开始有异步任务、消息投递和任务调度时,各个组件应该怎么分工。
一、先区分 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 或 RabbitMQworker:真正执行任务的进程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 + Beat。APScheduler 往往更轻量。
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 入门》:数据库访问是另一条并行主线