FastStream × FastAPI 集成实战:用 StreamRouter 将 Kafka/RabbitMQ/NATS 消息处理无缝接入 FastAPI 应用
【免费下载链接】faststreamAsynchronous Python framework for event-driven services. A thin client for Kafka, RabbitMQ, NATS, Redis and MQTT with full access to native broker features, plus AsyncAPI docs, in-memory tests and observability out of the box.项目地址: https://gitcode.com/GitHub_Trending/fa/faststream
本文以 docs/docs/en/getting-started/integrations/fastapi/index.md 为骨架,结合 faststream/_internal/fastapi/router.py 与 faststream/_internal/fastapi/route.py 等源码实现,系统讲解FastStream的FastAPI 插件:如何在 FastAPI 应用中直接声明消息订阅、复用 FastAPI 原生依赖注入与后台任务、通过四种方式访问 Broker、挂载 AsyncAPI 文档、编写内存测试以及拆分多个 Router。读完本文,你将能在一个 FastAPI 应用中同时承载 HTTP 接口与事件驱动消息处理,并完成测试与文档配置。
版本现状与迁移提示
在动手之前必须先明确一点:该集成已被官方标记为废弃(deprecated)。文档原文顶部与 faststream/kafka/fastapi/init.py 中的DeprecationWarning均明确说明:
The integration has been moved to the
faststream_fastapipackage and will be removed in1.0.0version.
pip install faststream_fastapi集成代码已被迁移到独立的faststream_fastapi包,并将在 FastStream 1.0.0 版本中从当前仓库移除。因此,本文介绍的 API 形态(StreamRouter、faststream.[broker].fastapi.Context等)在迁移后的包中依然延续,但新项目建议直接安装faststream_fastapi使用;本文同时记录当前仓库中的实现细节,供现有项目升级参考。
快速开始:在 FastAPI 中声明消息处理器
FastStream 可以作为 FastAPI 的一部分使用:只需要导入所需的StreamRouter(如KafkaRouter、RabbitRouter、NatsRouter、RedisRouter、MQTTRouter、ConfluentRouter),然后像声明普通 FastStream 应用一样声明消息处理器即可。
以 AIOKafka 为例,完整代码见 docs/docs_src/integrations/fastapi/kafka/base.py:
from fastapi import Depends, FastAPI from pydantic import BaseModel from faststream.kafka.fastapi import KafkaRouter, Logger router = KafkaRouter("localhost:9092") class Incoming(BaseModel): m: dict def call() -> bool: return True @router.subscriber("test") @router.publisher("response") async def hello(message: Incoming, logger: Logger, dependency: bool = Depends(call)): logger.info("Incoming value: %s, depends value: %s" % (message.m, dependency)) return {"response": "Hello, Kafka!"} @router.get("/") async def hello_http(): return "Hello, HTTP!" app = FastAPI() app.include_router(router)这段代码同时做到了三件事:
@router.subscriber("test")订阅 Kafka 的test主题;@router.publisher("response")将处理器返回值发布到response主题(发布/订阅链式声明与普通 FastStream 完全一致);@router.get("/")声明一个普通 HTTP 路由——这正是下文会讲到的StreamRouter 同时是 HttpRouter的能力。
各消息队列对应的导入路径与路由类完全同构:
| Broker | StreamRouter 导入 | 文档示例 |
|---|---|---|
| AIOKafka | from faststream.kafka.fastapi import KafkaRouter | docs/docs_src/integrations/fastapi/kafka/base.py |
| Confluent | from faststream.confluent.fastapi import KafkaRouter | docs/docs_src/integrations/fastapi/confluent/base.py |
| RabbitMQ | from faststream.rabbit.fastapi import RabbitRouter | docs/docs_src/integrations/fastapi/rabbit/base.py |
| NATS | from faststream.nats.fastapi import NatsRouter | docs/docs_src/integrations/fastapi/nats/base.py |
| Redis | from faststream.redis.fastapi import RedisRouter | docs/docs_src/integrations/fastapi/redis/base.py |
| MQTT | from faststream.mqtt.fastapi import MQTTRouter | docs/docs_src/integrations/fastapi/mqtt/base.py |
依赖注入体系切换:用 FastAPI 的原生能力
以这种模式使用时,FastStream 不再使用自己的依赖系统,而是整体融入 FastAPI。这意味着:
- 可以使用
fastapi.Depends、fastapi.BackgroundTasks等所有原生 FastAPI 特性,如同处理普通 HTTP 端点; - 不能使用
faststream.Context和faststream.Depends; - 需要上下文取值时,应使用
faststream.[broker].fastapi.Context(以及仓库预先创建好的注解别名,参见 docs/docs/en/getting-started/context.md 中的 Annotated Aliases 一节)。
上面的示例中dependency: bool = Depends(call)用的是fastapi.Depends,而非faststream.Depends——这是集成模式下最常见的坑。
源码层的强制校验
这种限制不是文档口头约定,而是在源码层做了强校验。faststream/_internal/fastapi/route.py 的wrap_callable_to_fastapi_compatible会在注册处理器时检查函数签名:
- 若发现
faststream.Depends(Dependant类型),直接抛出SetupError,提示改为fastapi.Depends; - 若发现
faststream.Context,同样抛出SetupError,提示改用faststream.[broker].fastapi.Context。
可见,混用两套依赖系统会在应用启动时就被拦截,而不是在运行时悄悄出错。
消息与请求参数的映射关系
处理 broker 消息时,整个消息体会被同时放入body与path两个请求参数,你可以按任意方便的方式访问它们;消息头则放入headers。这一定义的底层实现见 faststream/_internal/fastapi/route.py 中的StreamMessage类:
class StreamMessage(Request): def __init__(self, *, body, headers, path) -> None: self._headers = headers self._body = body self._query_params = path ...从源码看,消息体的分派规则是(route.py):
- 消息体是
dict时,body与path都取该字典; - 消息体是
list时,body取列表、path为空; - 其他标量类型时,
body与path均包装为{第一个参数名: 消息体},使消息可直接绑定到处理器签名中的首个位置参数。
因此你在处理器里既可以写def handler(message: Incoming)(经 Pydantic 校验),也可以把消息字段当作路径参数解构,两者底层数据源相同。
同时充当 HttpRouter:混合声明 HTTP 路由
StreamRouter 继承自 FastAPI 的APIRouter(源码见 faststream/_internal/fastapi/router.py 中class StreamRouter(APIRouter, StartAbleApplication, Generic[MsgType])),因此它可以完全当作一个 HttpRouter 使用,随意声明get、post、put等任何 HTTP 方法。上例第 20 行的@router.get("/")就是典型用法。
这带来一个非常自然的开发体验:同一个 Router 对象既管消息订阅,又管 HTTP 路由,而二者的处理器可以共用同一套 FastAPI 依赖体系。当你需要「HTTP 触发 → 发消息到 MQ」或「MQ 消息 → 写数据库」这类混合链路时,一个 Router 就能完成组织。
生命周期管理:lifespan 与 setup_state
由于 StreamRouter 需要管理 broker 连接,它会接管 FastAPI 对象的 lifespan。底层实现在 faststream/_internal/fastapi/router.py 的_wrap_lifespan中:
- 在 lifespan 启动阶段调用
self._start_broker()建立 broker 连接; - 执行所有
@after_startup钩子(见下文); - 若启用
setup_state,将{"broker": self.broker, ...}暴露到应用 state; - 关闭阶段执行
on_broker_shutdown钩子并调用self.broker.stop()。
版本兼容注意
- fastapi < 0.112.2:需要手动设置 lifespan,即
FastAPI(lifespan=router.lifespan_context); - fastapi >= 0.112.2:无需手动设置;若仍手动指定 lifespan,源码会发出
RuntimeWarning(见 router.py)。
setup_state=False 场景
如果你的ASGI 服务器不支持在 lifespan 内安装 state,可以禁用该行为:
router = StreamRouter(..., setup_state=False)代价是:之后无法从应用state中访问 broker(但它仍可通过router.broker访问)。该参数默认值为True,对应 router.py 的setup_state: bool = True。
访问 Broker 的四种方式
需要向 MQ 发送消息时,必须拿到 broker 对象。官方推荐与备选方式共四种:
方式一:Context 注解(推荐)
通过 Context 特性直接注入,类型注解由各 broker 的 fastapi 模块预先定义(如 faststream/kafka/fastapi/init.py 中的KafkaBroker = Annotated[KB, Context("broker")]):
from faststream.kafka.fastapi import KafkaBroker @router.get("/") async def handler(broker: KafkaBroker): ...各 Broker 对应:
from faststream.confluent.fastapi import KafkaBroker # Confluent from faststream.rabbit.fastapi import RabbitBroker # RabbitMQ from faststream.nats.fastapi import NatsBroker # NATS from faststream.redis.fastapi import RedisBroker # Redis from faststream.mqtt.fastapi import MQTTBroker # MQTT方式二:router.broker 直接访问
每个 Router 内部都持有 broker 实例,需要向 MQ 发消息时可直接使用,见 docs/docs_src/integrations/fastapi/kafka/send.py:
from fastapi import FastAPI from faststream.kafka.fastapi import KafkaRouter router = KafkaRouter("localhost:9092") app = FastAPI() @router.get("/") async def hello_http(): await router.broker.publish("Hello, Kafka!", "test") return "Hello, HTTP!" app.include_router(router)方式三:Depends 注入
如果要在程序的不同位置使用 broker,可以通过Depends注入,见 docs/docs_src/integrations/fastapi/kafka/depends.py:
from fastapi import Depends, FastAPI from typing_extensions import Annotated from faststream.kafka import KafkaBroker, fastapi router = fastapi.KafkaRouter("localhost:9092") app = FastAPI() def broker(): return router.broker @router.get("/") async def hello_http(broker: Annotated[KafkaBroker, Depends(broker)]): await broker.publish("Hello, Kafka!", "test") return "Hello, HTTP!" app.include_router(router)方式四:从应用 state 读取
若未通过setup_state=False禁用,可以从 FastAPI 应用 state 中直接读取:
from fastapi import Request @app.get("/") def main(request: Request): broker = request.state.broker使用 @after_startup 钩子
FastStream 应用提供了@after_startup钩子,允许在 broker 连接建立后执行一些操作(管理 broker 对象、发送消息等)。该钩子在FastAPI StreamRouter 上同样可用,示例见 docs/docs_src/integrations/fastapi/kafka/startup.py:
from fastapi import FastAPI from faststream.kafka.fastapi import KafkaRouter router = KafkaRouter("localhost:9092") @router.subscriber("test") async def hello(msg: str): return {"response": "Hello, Kafka!"} @router.after_startup async def test(app: FastAPI): await router.broker.publish("Hello!", "test") app = FastAPI() app.include_router(router)从源码看(router.py),after_startup与on_broker_shutdown均支持同步/异步函数:after_startup钩子在 broker 启动后依次执行,其返回值会合并进 lifespan 暴露给应用 state;on_broker_shutdown钩子在 broker 停止前执行。这为「启动时初始化队列/声明交换机」「关闭前清理资源」等场景提供了标准挂载点。
AsyncAPI 文档自动挂载
将 FastStream 用作 FastAPI 路由时,框架会自动在你的应用中注册承载 AsyncAPI 文档的端点,默认参数如下:
from faststream.kafka.fastapi import KafkaRouter router = KafkaRouter( ..., schema_url="/asyncapi", # 默认值 include_in_schema=True, # 默认值 )其余 Broker 完全一致(RabbitRouter、NatsRouter、RedisRouter、MQTTRouter等)。这样你将获得三个与 AsyncAPI schema 交互的路由:
| 路由 | 作用 |
|---|---|
/asyncapi | 与 CLI 生成的页面 相同的交互式文档页面 |
/asyncapi.json | 下载 JSON 格式的 schema |
/asyncapi.yaml | 下载 YAML 格式的 schema |
底层实现细节
这三个端点的实现在 faststream/_internal/fastapi/router.py 的_asyncapi_router方法中:
- HTML 页面由
get_asyncapi_html渲染,支持sidebar、info、servers、operations、messages、schemas等展示开关; - JSON/YAML 由
self.schema.to_specification()序列化生成; - 此外,文档页面还注册了
POST {schema_url}/try端点(include_in_schema=False),配合TryItOutProcessor实现在线「试一试」发布消息——前提是 broker 关联了测试 Broker,否则该端点会被跳过; - 这些端点仅在
include_in_schema=True时注册,否则docs_router为None。
值得注意的联动:在 lifespan 启动时,源码会用 FastAPI 应用的title、description、version、contact、license_info覆盖 AsyncAPI schema 的元信息(router.py),保证文档页与应用信息一致。
测试:用 TestClient 验证 StreamRouter
测试你的 FastAPI StreamRouter 时,仍然可以借助 FastStream 的内存测试 Broker(TestClient)。示例见 docs/docs_src/integrations/fastapi/kafka/test.py:
import pytest from faststream.kafka import TestKafkaBroker, fastapi router = fastapi.KafkaRouter() @router.subscriber("test") async def handler(msg: str): ... @pytest.mark.asyncio async def test_router(): async with TestKafkaBroker(router.broker) as br: await br.publish("Hi!", "test") handler.mock.assert_called_once_with("Hi!")核心要点:
- 测试时用
fastapi.KafkaRouter()不传连接参数(测试模式无需真实 broker); - 用
TestKafkaBroker(router.broker)包装真实 router 的 broker; - 通过
br.publish(...)模拟消息入队,然后用handler.mock.assert_called_once_with(...)断言处理器被正确调用且收到预期参数。
其余 Broker 对应TestConfluentBroker、TestRabbitBroker、TestNatsBroker、TestRedisBroker、TestMQTTBroker,测试方式完全一致。
多 Router 组织:拆分消息处理逻辑
像普通APIRouter一样,你仍然可以用多个 Router 拆分不同业务的消息处理逻辑。但这里有一个容易混淆的点:StreamRouter 会接管 FastAPI 对象的 lifespan,所以多个 StreamRouter 直接叠加可能会互相干扰。
官方推荐的方案是:使用普通的 FastStream Router(BrokerRouter)进行业务拆分,再 include 进 FastAPI 集成的那个 StreamRouter,就像 include 普通 broker 路由一样。这样还能方便地在 FastAPI 集成模式与纯 FastStream 应用之间复用端点。示例见 docs/docs_src/integrations/fastapi/kafka/router.py:
from fastapi import FastAPI from faststream.kafka import KafkaRouter from faststream.kafka.fastapi import KafkaRouter as StreamRouter core_router = StreamRouter() nested_router = KafkaRouter() @core_router.subscriber("core-topic") async def handler(): ... @nested_router.subscriber("nested-topic") async def nested_handler(): ... core_router.include_router(nested_router) app = FastAPI() app.include_router(core_router)源码中的约束
从 router.py 的include_router实现可以看出明确的约束:
include_router接受StreamRouter或BrokerRouter两种类型;- 传入普通 BrokerRouter 时,会为其中每个订阅者套上 FastAPI 兼容包装器(
_subscriber_compatibility_wrapper),再挂载到内部 broker——这是「复用端点」的关键机制; - 将一个 StreamRouter include 进另一个 StreamRouter 会直接抛出
TypeError,源码中的提示信息明确建议:用普通 broker 路由(如KafkaRouter、RabbitRouter)组织订阅者,再 include 进 StreamRouter,否则可能引发消息依赖返回EmptyPlaceholder等微妙的上下文问题。
小结与更多阅读
FastStream 的 FastAPI 插件让「HTTP + 事件驱动」两种编程模型在同一个应用中自然融合:消息处理器直接享受 FastAPI 的依赖注入、后台任务与测试体系,同时保留 FastStream 的发布/订阅链式声明、AsyncAPI 文档与内存测试能力。需要记住的要点:
- 依赖体系二选一:集成模式下用
fastapi.Depends与faststream.[broker].fastapi.Context,混用会在注册期被SetupError拦截; - 生命周期由 Router 接管:
fastapi >= 0.112.2无需手动设置 lifespan,特殊 ASGI 环境可setup_state=False; - Broker 访问四选一:Context 注解(推荐)、
router.broker、Depends、request.state.broker; - 文档开箱即用:
/asyncapi、/asyncapi.json、/asyncapi.yaml自动注册,schema 元信息与应用信息联动; - 多 Router 用 BrokerRouter 拆分,StreamRouter 之间不可互相 include。
若想深入了解依赖注入与上下文机制,可继续阅读 docs/docs/en/getting-started/context.md(Context 字段声明与注解别名);集成示例源码集中在 docs/docs_src/integrations/fastapi/ 目录,每个 Broker 均包含base.py、send.py、depends.py、startup.py、test.py、router.py六份可运行示例;对应的集成测试见 tests/docs/。注意,该集成将在 1.0.0 版本迁移至faststream_fastapi包,新项目请直接pip install faststream_fastapi。
【免费下载链接】faststreamAsynchronous Python framework for event-driven services. A thin client for Kafka, RabbitMQ, NATS, Redis and MQTT with full access to native broker features, plus AsyncAPI docs, in-memory tests and observability out of the box.项目地址: https://gitcode.com/GitHub_Trending/fa/faststream
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考