news 2026/9/17 18:54:12

FastStream × FastAPI 集成实战:用 StreamRouter 将 Kafka/RabbitMQ/NATS 消息处理无缝接入 FastAPI 应用

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
FastStream × FastAPI 集成实战:用 StreamRouter 将 Kafka/RabbitMQ/NATS 消息处理无缝接入 FastAPI 应用

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 等源码实现,系统讲解FastStreamFastAPI 插件:如何在 FastAPI 应用中直接声明消息订阅、复用 FastAPI 原生依赖注入与后台任务、通过四种方式访问 Broker、挂载 AsyncAPI 文档、编写内存测试以及拆分多个 Router。读完本文,你将能在一个 FastAPI 应用中同时承载 HTTP 接口与事件驱动消息处理,并完成测试与文档配置。

版本现状与迁移提示

在动手之前必须先明确一点:该集成已被官方标记为废弃(deprecated)。文档原文顶部与 faststream/kafka/fastapi/init.py 中的DeprecationWarning均明确说明:

The integration has been moved to thefaststream_fastapipackage and will be removed in1.0.0version.

pip install faststream_fastapi

集成代码已被迁移到独立的faststream_fastapi包,并将在 FastStream 1.0.0 版本中从当前仓库移除。因此,本文介绍的 API 形态(StreamRouterfaststream.[broker].fastapi.Context等)在迁移后的包中依然延续,但新项目建议直接安装faststream_fastapi使用;本文同时记录当前仓库中的实现细节,供现有项目升级参考。

快速开始:在 FastAPI 中声明消息处理器

FastStream 可以作为 FastAPI 的一部分使用:只需要导入所需的StreamRouter(如KafkaRouterRabbitRouterNatsRouterRedisRouterMQTTRouterConfluentRouter),然后像声明普通 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)

这段代码同时做到了三件事:

  1. @router.subscriber("test")订阅 Kafka 的test主题;
  2. @router.publisher("response")将处理器返回值发布到response主题(发布/订阅链式声明与普通 FastStream 完全一致);
  3. @router.get("/")声明一个普通 HTTP 路由——这正是下文会讲到的StreamRouter 同时是 HttpRouter的能力。

各消息队列对应的导入路径与路由类完全同构:

BrokerStreamRouter 导入文档示例
AIOKafkafrom faststream.kafka.fastapi import KafkaRouterdocs/docs_src/integrations/fastapi/kafka/base.py
Confluentfrom faststream.confluent.fastapi import KafkaRouterdocs/docs_src/integrations/fastapi/confluent/base.py
RabbitMQfrom faststream.rabbit.fastapi import RabbitRouterdocs/docs_src/integrations/fastapi/rabbit/base.py
NATSfrom faststream.nats.fastapi import NatsRouterdocs/docs_src/integrations/fastapi/nats/base.py
Redisfrom faststream.redis.fastapi import RedisRouterdocs/docs_src/integrations/fastapi/redis/base.py
MQTTfrom faststream.mqtt.fastapi import MQTTRouterdocs/docs_src/integrations/fastapi/mqtt/base.py

依赖注入体系切换:用 FastAPI 的原生能力

以这种模式使用时,FastStream 不再使用自己的依赖系统,而是整体融入 FastAPI。这意味着:

  • 可以使用fastapi.Dependsfastapi.BackgroundTasks等所有原生 FastAPI 特性,如同处理普通 HTTP 端点;
  • 不能使用faststream.Contextfaststream.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.DependsDependant类型),直接抛出SetupError,提示改为fastapi.Depends
  • 若发现faststream.Context,同样抛出SetupError,提示改用faststream.[broker].fastapi.Context

可见,混用两套依赖系统会在应用启动时就被拦截,而不是在运行时悄悄出错。

消息与请求参数的映射关系

处理 broker 消息时,整个消息体会被同时放入bodypath两个请求参数,你可以按任意方便的方式访问它们;消息头则放入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时,bodypath都取该字典;
  • 消息体是list时,body取列表、path为空;
  • 其他标量类型时,bodypath均包装为{第一个参数名: 消息体},使消息可直接绑定到处理器签名中的首个位置参数。

因此你在处理器里既可以写def handler(message: Incoming)(经 Pydantic 校验),也可以把消息字段当作路径参数解构,两者底层数据源相同。

同时充当 HttpRouter:混合声明 HTTP 路由

StreamRouter 继承自 FastAPI 的APIRouter(源码见 faststream/_internal/fastapi/router.py 中class StreamRouter(APIRouter, StartAbleApplication, Generic[MsgType])),因此它可以完全当作一个 HttpRouter 使用,随意声明getpostput等任何 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中:

  1. 在 lifespan 启动阶段调用self._start_broker()建立 broker 连接;
  2. 执行所有@after_startup钩子(见下文);
  3. 若启用setup_state,将{"broker": self.broker, ...}暴露到应用 state;
  4. 关闭阶段执行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_startupon_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 完全一致(RabbitRouterNatsRouterRedisRouterMQTTRouter等)。这样你将获得三个与 AsyncAPI schema 交互的路由:

路由作用
/asyncapi与 CLI 生成的页面 相同的交互式文档页面
/asyncapi.json下载 JSON 格式的 schema
/asyncapi.yaml下载 YAML 格式的 schema

底层实现细节

这三个端点的实现在 faststream/_internal/fastapi/router.py 的_asyncapi_router方法中:

  • HTML 页面由get_asyncapi_html渲染,支持sidebarinfoserversoperationsmessagesschemas等展示开关;
  • JSON/YAML 由self.schema.to_specification()序列化生成;
  • 此外,文档页面还注册了POST {schema_url}/try端点(include_in_schema=False),配合TryItOutProcessor实现在线「试一试」发布消息——前提是 broker 关联了测试 Broker,否则该端点会被跳过;
  • 这些端点仅在include_in_schema=True时注册,否则docs_routerNone

值得注意的联动:在 lifespan 启动时,源码会用 FastAPI 应用的titledescriptionversioncontactlicense_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 对应TestConfluentBrokerTestRabbitBrokerTestNatsBrokerTestRedisBrokerTestMQTTBroker,测试方式完全一致。

多 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接受StreamRouterBrokerRouter两种类型;
  • 传入普通 BrokerRouter 时,会为其中每个订阅者套上 FastAPI 兼容包装器(_subscriber_compatibility_wrapper),再挂载到内部 broker——这是「复用端点」的关键机制;
  • 将一个 StreamRouter include 进另一个 StreamRouter 会直接抛出TypeError,源码中的提示信息明确建议:用普通 broker 路由(如KafkaRouterRabbitRouter)组织订阅者,再 include 进 StreamRouter,否则可能引发消息依赖返回EmptyPlaceholder等微妙的上下文问题。

小结与更多阅读

FastStream 的 FastAPI 插件让「HTTP + 事件驱动」两种编程模型在同一个应用中自然融合:消息处理器直接享受 FastAPI 的依赖注入、后台任务与测试体系,同时保留 FastStream 的发布/订阅链式声明、AsyncAPI 文档与内存测试能力。需要记住的要点:

  1. 依赖体系二选一:集成模式下用fastapi.Dependsfaststream.[broker].fastapi.Context,混用会在注册期被SetupError拦截;
  2. 生命周期由 Router 接管fastapi >= 0.112.2无需手动设置 lifespan,特殊 ASGI 环境可setup_state=False
  3. Broker 访问四选一:Context 注解(推荐)、router.brokerDependsrequest.state.broker
  4. 文档开箱即用/asyncapi/asyncapi.json/asyncapi.yaml自动注册,schema 元信息与应用信息联动;
  5. 多 Router 用 BrokerRouter 拆分,StreamRouter 之间不可互相 include。

若想深入了解依赖注入与上下文机制,可继续阅读 docs/docs/en/getting-started/context.md(Context 字段声明与注解别名);集成示例源码集中在 docs/docs_src/integrations/fastapi/ 目录,每个 Broker 均包含base.pysend.pydepends.pystartup.pytest.pyrouter.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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 18:48:50

STM32 Bootloader与串口IAP实现:从原理到代码,轻松搞定固件升级

很多人刚接触单片机开发时&#xff0c;对 bootloader 总是有一种“高深莫测”的感觉。后来我自己做产品&#xff0c;才真正意识到它其实就是一段“先于主程序运行的小程序”&#xff0c;并没有想象中复杂。尤其是当你需要给已经出货的设备做固件升级时&#xff0c;IAP&#xff…

作者头像 李华
网站建设 2026/9/17 18:45:21

S7-1500 PLC硬件配置、电源预算与硬件诊断实战

简介&#xff1a;这份《S7-1500 PLC应用技术》第2章课件&#xff0c;面向自动化、电气控制专业学生及初学S7-1500的工程人员&#xff0c;用于梳理该系列PLC硬件体系与选型要点。内容分六部分&#xff1a;SIMATIC产品定位、CPU模块、电源模块、信号模块、通信模块与CPU操作模式。…

作者头像 李华
网站建设 2026/9/17 18:41:18

VS2015安装包损坏修复全指南:校验、离线重建与工具链提取

1. 这不是“重装就完事”的问题&#xff1a;VS2015安装包损坏/丢失的真实战场你点开那个下载了三小时的 vs2015community.exe&#xff0c;双击后弹出“无法验证安装包完整性”、“找不到 bootstrapper.exe”、“setup.exe 已损坏”——不是你的网速慢&#xff0c;也不是磁盘满了…

作者头像 李华