06 - FastAPI + RabbitMQ 的整合架构思路
前面都在讲原理,这一章落到你自己的项目:RabbitMQ 接进 FastAPI 后,整个系统长什么样、各部分怎么分工。 本章重架构视角,代码只给”意思到了”的骨架,帮你建立整体图景。
6.1 整体架构:三个独立的角色
引入 RabbitMQ 后,你的系统从”一个 FastAPI 应用”变成三块独立部署的东西:
用户请求
│
▼
┌─────────────────────┐ ┌──────────────────────┐
│ FastAPI 应用 │ │ Worker 消费者 │
│ (生产者) │ │ (独立进程/容器) │
│ │ │ │
│ 接口里: │ │ 循环从队列取消息: │
│ 1. 处理核心逻辑 │ │ 1. 收到 user_registered│
│ 2. 写数据库 │ │ 2. 发邮件/短信 │
│ 3. 发一条消息 ───────┼──┐ ┌──┼─▶3. 处理成功 → ACK │
│ 4. 立刻返回响应 │ │ │ │ │
└─────────────────────┘ │ │ └──────────────────────┘
▼ │
┌───────────────────────┐
│ RabbitMQ Broker │
│ (独立服务 / 容器) │
│ 交换机 → 队列 → 暂存 │
└───────────────────────┘
▲
┌─────┴──────┐
│ 数据库 PG │ ← 生产者和消费者按需各自读写
│ 缓存 Redis │
└────────────┘
三个角色各司其职:
| 角色 | 是什么 | 怎么运行 |
|---|---|---|
| FastAPI 应用(生产者) | 你现在写的 Web 服务 | uvicorn 跑起来,处理 HTTP 请求,在接口里发消息 |
| RabbitMQ Broker | 消息中间件 | 独立服务,用 Docker 单独跑一个容器 |
| Worker(消费者) | 一个独立的 Python 脚本 | 单独用 python worker.py 跑,是一个独立进程,不在 FastAPI 里 |
最关键的架构认知:消费者 Worker 是独立于 FastAPI 单独运行的进程。它不是一个 API 接口,而是一个”一直挂着、循环从队列取消息处理”的常驻程序。这是初学者最容易迷糊的点。
6.2 常用的 Python 库
| 库 | 说明 | 适合 |
|---|---|---|
| pika | 最经典的 RabbitMQ Python 客户端,同步风格,概念直白 | 学习原理、写独立 worker(推荐入门用它理解) |
| aio-pika | 基于 asyncio 的异步客户端 | 想和 FastAPI 的 async 风格统一时 |
| Celery | 上层任务框架,把发任务/执行任务封装得很好用 | 生产项目里快速落地异步任务 |
学习阶段建议先用 pika 把生产者/消费者手写一遍,彻底理解消息怎么流动;理解之后,真实项目再上 Celery 提效。
6.3 生产者:在 FastAPI 接口里发消息
生产者的职责很轻——连上 Broker,把消息往交换机一扔,就返回。骨架大致是这样(示意,不用抠细节):
import pika, json
from fastapi import FastAPI
app = FastAPI()
@app.post("/register")
def register(user: dict):
# 1. 核心逻辑:写数据库(必须同步做完)
user_id = save_user_to_db(user)
# 2. 发一条"用户已注册"消息给 RabbitMQ,然后就不管了
conn = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
channel = conn.channel()
channel.basic_publish(
exchange="user_events",
routing_key="user.registered",
body=json.dumps({"user_id": user_id}),
)
conn.close()
# 3. 立刻返回,不等邮件/短信
return {"status": "ok", "user_id": user_id}
注意这里的思路:接口只负责”写库 + 发消息”,发邮件发短信这些一个字都没写——那是消费者的事。这就是”解耦”落到代码上的样子。
(工程上不会每次都新建连接,而是复用连接/连接池,这里为了讲清楚流程做了简化。)
6.4 消费者:一个独立常驻的 worker
消费者是单独一个 worker.py,它连上 Broker、盯着队列、来一条处理一条:
import pika, json
def on_message(channel, method, properties, body):
data = json.loads(body)
send_welcome_email(data["user_id"]) # 真正干活
channel.basic_ack(delivery_tag=method.delivery_tag) # 处理成功 → 确认
conn = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
channel = conn.channel()
channel.queue_declare(queue="email_queue", durable=True)
channel.basic_consume(queue="email_queue", on_message_callback=on_message)
print("Worker 已启动,等待消息…")
channel.start_consuming() # 常驻循环,一直挂着
关键理解:
- 它不是 FastAPI 的一部分,用
python worker.py单独启动。 start_consuming()会让它一直挂着,有消息就处理,没消息就等着。- 处理成功后
basic_ack确认,Broker 才删消息(第 07 章 细讲)。 - 想处理更快?多起几个 worker 进程,它们自动分摊队列里的消息(对应 第 02 章 的点对点负载均衡)。
6.5 用 Docker 把整套跑起来(架构视角)
你已经学过 Docker,正好用它把这套架构编排起来。概念上是这样:
docker-compose 里定义几个服务:
┌────────────────────────────────────────────────┐
│ rabbitmq :官方镜像 rabbitmq:3-management │
│ 暴露 5672(通信) + 15672(管理后台) │
│ │
│ api :你的 FastAPI 镜像 │
│ 依赖 rabbitmq,连它发消息 │
│ │
│ worker :同一套代码,启动命令改成 │
│ python worker.py(跑消费者) │
│ 可以起多个副本分摊任务 │
│ │
│ db / redis :你原有的 PostgreSQL / Redis │
└────────────────────────────────────────────────┘
注意:api 和 worker 往往是同一份代码、同一个镜像,只是启动命令不同——一个跑
uvicorn(当生产者),一个跑worker.py(当消费者)。这是很常见的实践。
启动后,浏览器打开 http://localhost:15672(默认账号 guest/guest)就能看到 RabbitMQ 的管理后台,直观地看到队列里有多少消息、消费速度多快、有没有堆积。
6.6 消息流动全景(把本章串起来)
用户 POST /register
│
▼
[FastAPI] 写库 → 发消息 "user.registered" → 立刻返回"注册成功"(0.05s)
│ │
│ ▼
│ [RabbitMQ 交换机]
│ 按路由分发
│ ┌──────────┴──────────┐
│ ▼ ▼
│ [email_queue] [stats_queue]
│ │ │
│ ▼ ▼
│ [邮件 worker] [统计 worker]
│ 发欢迎邮件 更新统计
│ │ │
│ ACK 确认 ACK 确认
▼
用户早已拿到响应,后台任务在慢慢跑
本章小结
| 知识点 | 要点 |
|---|---|
| 三个角色 | FastAPI(生产者)、RabbitMQ(Broker)、Worker(消费者),独立部署 |
| 最易错点 | 消费者是独立常驻进程(python worker.py),不是 API 接口 |
| 生产者职责 | 接口里只做”核心逻辑 + 发消息”,然后立刻返回 |
| 消费者职责 | 循环取消息 → 处理 → ACK 确认;多起几个进程即可扩容 |
| 常用库 | 入门用 pika 理解原理,生产用 Celery 提效 |
| 部署 | 用 Docker Compose 编排;api 和 worker 常共用一个镜像、启动命令不同 |
下一章预告:最后聊聊工程上绕不开的可靠性话题——消息会丢吗?会重复吗?处理失败怎么办?(只讲思路,让你心里有数。)
上一章 ← 05 - RabbitMQ 的架构与工作原理 | 下一章 → 07 - 可靠性与常见坑