FastStream Redis List 发布实战:用 `@broker.publisher(list=...)` 构建 List 消息管道 📅 发布时间:2026/9/18 13:22:49 👁 浏览次数: FastStream Redis List 发布实战用broker.publisher(list...)构建 List 消息管道【免费下载链接】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/faststreamFastStream 是面向事件驱动服务的异步 Python 框架其 Redis 模块原生支持将 Redis List 当作消息队列使用既能用broker.subscriber(list...)从 List 中阻塞弹出BLPOP/LPOP消费消息也能用broker.publisher(list...)把处理结果推回另一个 List从而用两个装饰器组合出一条“输入 List → 处理逻辑 → 输出 List”的流水线。本文以 FastStream 仓库的官方示例list_pub.py为主线结合 Redis 发布器实现 与 List 订阅器实现 等源码讲解如何在 FastStream 中完成 List 的发布/订阅以及底层是如何调用 Redis 命令的。理解 Redis List 发布机制与 Redis Streams 类似消息同样可以被发布到 Redis List。FastStream 通过#!python broker.publisher(...)装饰器配合 List 名称将消息推送到指定的 List 上。整个模型非常简单List 在 Redis 中是元素有序的队列结构天然适合“先进先出”的任务队列场景FastStream 的发布器不直接面向PUBLISH/SUBSCRIBE的发布订阅语义而是通过LPUSH/RPUSH之类的写操作把消息写进 List订阅端则通过阻塞弹出BLPOP或非阻塞弹出LPOP把消息取出形成完整的队列式收发闭环。从源码上看发布器在 factory.py 中根据channel、list、stream三个目标参数二选一创建对应实现当list参数被传入时会实例化ListPublisher批量模式则实例化ListBatchPublisher。这保证了同一个broker.publisher装饰器既能面向 Channel也能面向 List 或 Stream。分步构建从订阅到发布下面按官方示例的步骤完整展示如何在 FastStream 中搭建一个基于 Redis List 的发布管道。1. 实例化 RedisBrokerbroker RedisBroker(redis://localhost:6379)RedisBroker是 FastStream Redis 模块的入口构造参数为 Redis 连接地址默认redis://localhost:6379。Broker 负责管理连接、注册订阅者与发布器并生成对应规格specification与用例usecase对象。2. 创建 FastStream 应用app FastStream(broker)FastStream(broker)将 Broker 包装为完整的应用对象提供生命周期管理启动、优雅关闭以及可选的 AsyncAPI 文档生成能力。3. 定义 Pydantic 数据模型class Data(BaseModel): data: NonNegativeFloat Field( ..., examples[0.5], descriptionFloat data example, )示例使用 Pydantic v2 的BaseModel定义消息体。字段data是NonNegativeFloat非负浮点数Field(...)表示必填examples与description除供文档展示外也会被 AsyncAPI 规范拾取用于生成消息 schema。FastStream 在消费端会按该模型自动做反序列化与校验发布端则会把函数返回值序列化回字节流写入 Redis。4. 用装饰器组合出“订阅→处理→发布”函数broker.subscriber(listinput-list) broker.publisher(listoutput-list) async def on_input_data(msg: Data) - Data: return Data(datamsg.data 1.0)这是整个示例的核心模式broker.subscriber(listinput-list)让该函数成为input-list的消费者从该 List 取出消息并反序列化为Databroker.publisher(listoutput-list)函数返回值会自动被序列化并推送到output-list函数体内完成业务逻辑这里是将msg.data加 1.0。装饰器叠加后形成一条数据管道从input-list读取消息 → 应用逻辑 → 把结果写入output-list。这种写法让“输入 List / 输出 List”的关系在声明处一目了然无需手工管理发布器对象。完整示例代码以下是官方示例的完整代码来源list_pub.pyfrom pydantic import BaseModel, Field, NonNegativeFloat from faststream import FastStream from faststream.redis import RedisBroker class Data(BaseModel): data: NonNegativeFloat Field( ..., examples[0.5], descriptionFloat data example, ) broker RedisBroker(redis://localhost:6379) app FastStream(broker) broker.subscriber(listinput-list) broker.publisher(listoutput-list) async def on_input_data(msg: Data) - Data: return Data(datamsg.data 1.0)整个流程中消息从输入 List 出队dequeue、被处理、再入队enqueue到输出 List开发者可以借助 Redis 这种快速的内存数据结构把 List 当作轻量级消息队列来使用。深入源码List 发布器如何工作发布器的工厂选择与校验在 factory.py 中create_publisher首先调用validate_options(channel..., list..., stream...)校验目标参数确保三者中至少指定一个且不冲突随后根据传入的list参数构造ListPublisher或ListBatchPublisher由ListSub.batch决定。ListSub定义于 list_sub.py其参数包括list_nameList 名称batch是否批量发布/消费max_records批量模式下的最大记录数默认 10polling_interval轮询间隔默认 0.1 秒。ListPublisher 的发布过程ListPublisher 的publish方法构造一个RedisPublishCommand见 response.py该命令通过set_destination(list...)将目标类型标记为DestinationType.List随后交给底层 producer 执行对应 Redis 写命令。发布器还支持在声明时通过headers、reply_to等参数附加元数据config.py使用request方法实现 List 上的请求-响应模式内部携带timeout默认 30 秒使用pipeline参数把发布操作合并进 Redis 管道Pipeline批量执行。值得注意的细节ListPublisher通过self._list.add_prefix(self._outer_config.prefix)支持全局前缀prefix自动拼接也就是说 Broker 级prefix配置会对 List 名称统一加前缀方便多环境/多租户隔离。批量发布ListBatchPublisher若在声明时使用批量模式batchTrue会生成 ListBatchPublisher。它的publish接收可变数量的消息*messages一次调用即可把多条消息写入 List其内部通过_basic_publish_batch执行批量写入。源码注释还特别说明见_publish实现对应 issue #3056批量发布器返回空结果时会退化为写入一个空消息而非“零条消息的批次”以保证与 Redis 命令语义一致。消费端视角List 订阅器如何取消息发布只是管道的一半FastStream 的 ListSubscriber 定义了消费逻辑单条模式_get_msgs使用client.blpop(list_name, timeoutpolling_interval)阻塞弹出单条消息阻塞时长由polling_interval默认 0.1s控制批量模式ListBatchSubscriber 使用client.lpop(name..., countmax_records)一次性弹出最多max_records默认 10条消息构成BatchListMessage交给处理器若没有取到消息则按polling_interval休眠后继续轮询并发模式ListConcurrentSubscriber叠加ConcurrentMixin支持并发消费。此外订阅器还提供无处理器模式下的get_one(timeout...)与__aiter__迭代接口list_subscriber.py方便在测试或脚本中手工取消息。如何测试 List 发布管道FastStream 提供内存测试能力无需真实 Redis 即可验证发布逻辑。官方测试位于 tests/docs/redis/list/test_list_pub.pyimport pytest from faststream.redis import TestRedisBroker pytest.mark.redis() pytest.mark.asyncio() async def test_list_publisher() - None: from docs.docs_src.redis.list.list_pub import broker, on_input_data publisher list(broker.publishers)[0] # noqa: RUF015 async with TestRedisBroker(broker) as br: await br.publish({data: 1.0}, listinput-list) on_input_data.mock.assert_called_once_with({data: 1.0}) publisher.mock.assert_called_once_with({data: 2.0})测试要点用TestRedisBroker(broker)以 in-memory 模式启动 Broker不连接任何真实 Redis通过br.publish({data: 1.0}, listinput-list)向输入 List 注入消息断言订阅函数on_input_data.mock被调用且入参为{data: 1.0}断言发布器publisher.mock被调用且入参为{data: 2.0}即处理结果1.0 1.0被正确发布到输出 List。这种 mock 机制让 List 管道的单元测试既快速又稳定是 CI 中验证消息流逻辑的推荐方式。与其他 Redis 目标Channel / Stream的对比FastStream 的broker.publisher同样支持 Channel 与 Stream 目标三者底层命令与适用场景不同依据 response.py 的DestinationType枚举与各 usecase 实现目标类型声明方式底层语义典型场景Listpublisher(list...)写入 Redis ListLPUSH/RPUSH等消费端BLPOP/LPOP任务队列、worker 分发、流水线中间产物Channelpublisher(channel...)Redis Pub/SubPUBLISH实时广播实时通知、扇出广播Streampublisher(stream...)Redis Stream支持maxlen等事件溯源、日志、可靠消息重放选择 List 而不是 Stream 时主要看重其“队列”语义的简洁与低开销若需要持久化、消费者组、消息确认等能力则应转向 Stream。FastStream 对三者提供一致的装饰器式声明语法迁移成本很低。小结通过broker.subscriber(list...)与broker.publisher(list...)两个装饰器的组合FastStream 让开发者可以用几行代码把 Redis List 变成一条可运行的异步消息管道。其底层由 ListPublisher 负责写入、ListSubscriber 负责阻塞弹出消费并通过TestRedisBroker支持零依赖的内存测试。如果你想快速体验可参考 examples/redis/list_sub.py订阅侧示例与官方文档 publishing.md本文依据的原始文档结合本仓库的 list_pub.py 搭建你自己的 List 消息管道。【免费下载链接】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),仅供参考