首页 / 知识库 / python服务端进阶 / 消息队列入门

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 - 可靠性与常见坑