
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 thefaststream_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.pyfrom 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的能力。各消息队列对应的导入路径与路由类完全同构BrokerStreamRouter 导入文档示例AIOKafkafrom faststream.kafka.fastapi import KafkaRouterdocs/docs_src/integrations/fastapi/kafka/base.pyConfluentfrom faststream.confluent.fastapi import KafkaRouterdocs/docs_src/integrations/fastapi/confluent/base.pyRabbitMQfrom faststream.rabbit.fastapi import RabbitRouterdocs/docs_src/integrations/fastapi/rabbit/base.pyNATSfrom faststream.nats.fastapi import NatsRouterdocs/docs_src/integrations/fastapi/nats/base.pyRedisfrom faststream.redis.fastapi import RedisRouterdocs/docs_src/integrations/fastapi/redis/base.pyMQTTfrom faststream.mqtt.fastapi import MQTTRouterdocs/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.DependsDependant类型直接抛出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(lifespanrouter.lifespan_context)fastapi 0.112.2无需手动设置若仍手动指定 lifespan源码会发出RuntimeWarning见 router.py。setup_stateFalse 场景如果你的ASGI 服务器不支持在 lifespan 内安装 state可以禁用该行为router StreamRouter(..., setup_stateFalse)代价是之后无法从应用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.pyfrom 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.pyfrom 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_stateFalse禁用可以从 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.pyfrom 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.pyafter_startup与on_broker_shutdown均支持同步/异步函数after_startup钩子在 broker 启动后依次执行其返回值会合并进 lifespan 暴露给应用 stateon_broker_shutdown钩子在 broker 停止前执行。这为「启动时初始化队列/声明交换机」「关闭前清理资源」等场景提供了标准挂载点。AsyncAPI 文档自动挂载将 FastStream 用作 FastAPI 路由时框架会自动在你的应用中注册承载 AsyncAPI 文档的端点默认参数如下from faststream.kafka.fastapi import KafkaRouter router KafkaRouter( ..., schema_url/asyncapi, # 默认值 include_in_schemaTrue, # 默认值 )其余 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_schemaFalse配合TryItOutProcessor实现在线「试一试」发布消息——前提是 broker 关联了测试 Broker否则该端点会被跳过这些端点仅在include_in_schemaTrue时注册否则docs_router为None。值得注意的联动在 lifespan 启动时源码会用 FastAPI 应用的title、description、version、contact、license_info覆盖 AsyncAPI schema 的元信息router.py保证文档页与应用信息一致。测试用 TestClient 验证 StreamRouter测试你的 FastAPI StreamRouter 时仍然可以借助 FastStream 的内存测试 BrokerTestClient。示例见 docs/docs_src/integrations/fastapi/kafka/test.pyimport 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 RouterBrokerRouter进行业务拆分再 include 进 FastAPI 集成的那个 StreamRouter就像 include 普通 broker 路由一样。这样还能方便地在 FastAPI 集成模式与纯 FastStream 应用之间复用端点。示例见 docs/docs_src/integrations/fastapi/kafka/router.pyfrom 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_stateFalseBroker 访问四选一Context 注解推荐、router.broker、Depends、request.state.broker文档开箱即用/asyncapi、/asyncapi.json、/asyncapi.yaml自动注册schema 元信息与应用信息联动多 Router 用 BrokerRouter 拆分StreamRouter 之间不可互相 include。若想深入了解依赖注入与上下文机制可继续阅读 docs/docs/en/getting-started/context.mdContext 字段声明与注解别名集成示例源码集中在 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),仅供参考