尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

FastStream Confluent Kafka 基础订阅者实战指南:用 @broker.subscriber 消费单 Topic 与多 Topic 消息

发布时间:2026/9/18 5:05:38

资讯中心
01
ARTICLE

FastStream Confluent Kafka 基础订阅者实战指南:用 @broker.subscriber 消费单 Topic 与多 Topic 消息

FastStream Confluent Kafka 基础订阅者实战指南:用 @broker.subscriber 消费单 Topic 与多 Topic 消息
FastStream Confluent Kafka 基础订阅者实战指南用 broker.subscriber 消费单 Topic 与多 Topic 消息【免费下载链接】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 通过broker.subscriber(...)装饰器让开发者可以用声明式的方式订阅 Apache Kafka Topic并把消费函数自动注册为消息处理器。本文以 FastStream 官方文档的 Confluent Basic Subscriber 章节为骨架完整讲解如何导入KafkaBroker、用 Pydantic 定义消息结构、创建消费者函数并订阅单 Topic 与多 Topic同时结合仓库源码剖析订阅者背后的参数体系、ack 策略约束与三种订阅者实现类型帮助你写出可复制、可运行且行为可控的 Kafka 消费端代码。本文对应的官方文档位于 docs/docs/en/confluent/Subscriber/index.md配套可运行示例见 docs/docs_src/confluent/consumes_basics/app.py 与 docs/docs_src/confluent/multiple_topics_subscription/app.py。前提安装与 Broker 初始化本文使用基于 Confluent Kafka Python 客户端实现的KafkaBroker与aiokafka实现的 KafkaBroker 是两套独立实现见 docs/docs/en/confluent/index.md。安装命令为pip install faststream[confluent]0.4.0Confluent 支持从 FastStreamv0.4.0rc0版本开始提供。安装完成后在你的本机或测试环境启动 Kafka 服务默认地址localhost:9092即可开始消费。一、完整的订阅者应用Hello World 消息消费开始消费 Kafka Topic 非常简单把消费函数用broker.subscriber(...)装饰器修饰并以字符串形式传入 Topic 名称作为参数即可。下面是一个完整的 FastStream 应用它消费hello_worldTopic 上的HelloWorld消息from pydantic import BaseModel, Field from faststream import FastStream, Logger from faststream.confluent import KafkaBroker class HelloWorld(BaseModel): msg: str Field( ..., examples[Hello], descriptionDemo hello world message, ) broker KafkaBroker(localhost:9092) app FastStream(broker) broker.subscriber(hello_world) async def on_hello_world(msg: HelloWorld, logger: Logger): logger.info(msg)以上代码就是 docs/docs_src/confluent/consumes_basics/app.py 的完整内容。运行该应用后任何生产者向hello_worldTopic 发布消息on_hello_world都会被自动触发。下面按构建顺序拆解每一步。二、导入 FastStream 与 KafkaBroker要使用broker.subscriber(...)装饰器第一步是导入基础应用类FastStream和 Confluent 版的KafkaBroker并创建 broker 实例from faststream import FastStream, Logger from faststream.confluent import KafkaBroker其中FastStream是框架的应用入口负责启动/停止 broker、管理生命周期如after_startup钩子、生成 AsyncAPI 文档等KafkaBroker是 Confluent 版 Kafka 的 broker 类承担连接管理、订阅注册与发布能力Logger是 FastStream 注入的日志类型注解用于在消费函数中直接打印消息无需手动初始化 logger。三、用 Pydantic 定义 HelloWorld 消息结构接下来定义要从 Topic 中消费的消息结构。文档示例使用 Pydantic 模型字段用Field(...)标记为必填并附带examples与description元数据class HelloWorld(BaseModel): msg: str Field( ..., examples[Hello], descriptionDemo hello world message, )关键点说明消息体的 JSON 载荷会自动按该模型反序列化字段类型校验失败时会走框架的错误处理流程示例保持了简单的单字段结构但实际项目中完全可以定义任意复杂的嵌套 Pydantic 模型消息结构除了作为解析目标外还会被 FastStream 用于生成 AsyncAPI Schema订阅操作的 message payload 定义因此description等元数据会同步进入 API 文档。四、创建 KafkaBroker 并包装进 FastStreambroker KafkaBroker(localhost:9092) app FastStream(broker)KafkaBroker(localhost:9092)传入的是 Kafka broker 地址。将 broker 包装进FastStream后就可以用 FastStream 提供的 CLI 命令启动应用例如faststream run框架会负责 broker 的连接与消费任务的启动。五、编写消费函数并理解消息注入机制broker.subscriber(hello_world) async def on_hello_world(msg: HelloWorld, logger: Logger): logger.info(msg)这段代码的运行时行为如下被broker.subscriber(...)装饰的函数会被注册为hello_worldTopic 的处理器当消息被生产到 Kafka 的hello_worldTopic 时该函数会被调用消息被注入到类型化参数msg中参数的类型注解HelloWorld被用来解析消息体——框架读取消息字节流后按 JSON 反序列化并校验为HelloWorld实例最终on_hello_world会收到解析后的HelloWorld对象作为msg实参logger.info(msg)打印日志。这就是 FastStream 的核心消费模型类型注解即解析方案消费函数的签名决定了消息如何被解码无需手写任何 JSON 解析或类型转换代码。六、订阅多个 Topic一次订阅与堆叠装饰器的区别单个broker.subscriber(...)调用可以通过多个位置参数订阅多个 Topicfrom faststream import FastStream, Logger from faststream.confluent import KafkaBroker broker KafkaBroker(localhost:9092) app FastStream(broker) broker.subscriber(topic1, topic2, topic3) async def on_multiple_topics(msg: str, logger: Logger): logger.info(msg)对应的完整示例见 docs/docs_src/confluent/multiple_topics_subscription/app.py。这里有两个需要区分的概念一个订阅者一个 handler订阅多个 Topic如上面的写法所有列出的 Topictopic1、topic2、topic3的消息都会被同一个处理器接收且这些 Topic 属于同一个消费者组one consumer group。这意味着组内同一分区消息只会被该 handler 处理一次适合多个业务 Topic 汇聚到同一处理逻辑的场景堆叠多个broker.subscriber(...)装饰器对同一个函数使用多个 subscriber 装饰器会创建多个相互独立separate independent的处理器实例。每个处理器拥有自己的消费循环与生命周期即使订阅相同的 Topic也是完全独立的两套消费任务。在源码层面这两种方式的区别体现在 faststream/confluent/broker/registrator.py 的KafkaRegistrator.subscriber()实现中每次调用subscriber()都会通过create_subscriber()创建一个订阅者实例并注册到 broker 的处理器集合calls中多 Topic 是一个订阅者实例持有多个 topic堆叠装饰器则是多个订阅者实例。关于 max_workers 与 AckPolicy 的重要警告官方文档在此处给出了一条明确的约束max_workers 1只兼容AckPolicy.ACK_FIRST。如果以max_workers 1配合任何其他 ack policy无论订阅多少个 Topic都会在启动时抛出SetupError。也就是说当你希望通过max_workers开启并发消息处理时必须同时保持默认的AckPolicy.ACK_FIRST若想改用AckPolicy.ACK或AckPolicy.MANUAL等策略则不能与max_workers 1组合使用。这一约束在订阅者工厂创建时被校验具体逻辑可追溯 faststream/confluent/subscriber/factory.py。七、深入源码subscriber 的完整参数体系与三种订阅者类型broker.subscriber(...)并不是一个魔法语法它最终会走到KafkaRegistrator.subscriber()的完整签名见 faststream/confluent/broker/registrator.py。除了 Topic 名称之外你还可以传入以下常用参数参数默认值作用*topics必填Kafka Topic 名称可传多个字符串或Topic对象用Topic对象可配置 Topic 创建细节partitions()指定消费的具体分区TopicPartition序列不设置则使用消费者组动态分区分配polling_interval0.1轮询间隔秒group_idNone消费者组 ID为None时禁用组管理与自动提交 offsetgroup_instance_idNone消费者实例唯一标识设置后成为静态组员不参与组管理与再均衡fetch_max_wait_ms500服务端在未满足fetch_min_bytes时最长阻塞等待的时间毫秒fetch_max_bytes50 * 1024 * 1024单次 fetch 请求服务端可返回的数据上限字节注意消费端会对多个 broker 并行 fetchfetch_min_bytes1fetch 请求最少返回数据量不足则最多等待fetch_max_wait_msmax_partition_fetch_bytes1 * 1024 * 1024单分区最多返回数据量内存峰值约为分区数 × max_partition_fetch_bytesauto_offset_resetlatest无有效 offset 时的重置策略earliest/latest/noneauto_commit_interval_ms5 * 1000自动提交 offset 的间隔毫秒check_crcsTrue是否校验消息 CRC32追求极致性能时可关闭partition_assignment_strategy(roundrobin,)消费者组内分区分配策略列表max_poll_interval_ms5 * 60 * 1000两次批量消费之间的最大允许时间超时将被视为故障并触发再均衡session_timeout_ms10 * 1000组会话与故障检测超时超时未收到心跳会被移除并触发再均衡heartbeat_interval_ms3 * 1000心跳间隔通常应低于session_timeout_ms一般不超过其 1/3isolation_levelread_uncommitted事务消息读取级别read_committed只返回已提交事务消息read_uncommitted返回全部消息on_assign/on_revoke/on_lostNone再均衡回调分别在分区分配、撤销、丢失时被调用batchFalse是否以批量方式消费max_recordsNone批量消费时一次最多获取的消息条数max_workersNone并发处理消息的工作线程数大于 1 时启用并发模式ack_policyAckPolicy.ACK_FIRST消费确认策略详见下文persistentTrue订阅者是否常驻运行dependencies/parser/decoder/codec框架默认依赖注入、消息解析器、解码器与自定义编解码器title/description/include_in_schema框架默认AsyncAPI 规范文档相关配置在subscriber()实现内部faststream/confluent/broker/registrator.py会根据参数组合实例化三种不同的订阅者定义见 faststream/confluent/subscriber/usecase.pyDefaultSubscriber默认的单条消费模式batchFalse且max_workers1。其get_msg()通过consumer.getone(timeoutpolling_interval)逐条拉取消息BatchSubscriberbatchTrue时的批量消费模式。get_msg()调用consumer.getmany(timeout..., max_records...)消费函数收到的msg是消息元组详见 docs/docs/en/confluent/Subscriber/batch_subscriber.md示例见 docs/docs_src/confluent/batch_consuming_pydantic/app.pyConcurrentDefaultSubscribermax_workers 1时的并发模式继承ConcurrentMixin通过start_consume_task()启动并发消费任务。# 批量订阅示例msg 声明为 list 且 batchTrue broker.subscriber(test_batch, batchTrue) async def handle_batch(msg: list[HelloWorld], logger: Logger): logger.info(msg)八、ack 确认策略与订阅者强相关的行为配置虽然本文聚焦基础订阅但 ack 策略直接决定了消息被处理后 offset 何时提交是订阅者行为的关键一环完整文档见 docs/docs/en/confluent/ack.md。默认情况下 FastStream 使用AckPolicy.ACK_FIRST由 Kafka 客户端通过enable.auto.commit自动提交 offset即至多一次at most once消费语义。若需要至少一次at least once语义可设置消费者组并改用AckPolicy.ACK——offset 在 handler 执行成功后才提交异常时同样会提交broker.subscriber(test, group_idgroup, ack_policyAckPolicy.ACK) async def base_handler(body: str): ...若希望失败后重试可用AckPolicy.NACK_ON_ERROR——出错时 offset 不提交消费端 seek 回该消息重新读取broker.subscriber(test, group_idgroup, ack_policyAckPolicy.NACK_ON_ERROR) async def base_handler(body: str): ...各策略的行为差异如下表摘自 docs/docs/en/confluent/ack.mdAckPolicy成功时出错时说明MANUAL不做任何事不做任何事消费者从不提交 offset完全手动控制ACK_FIRST不做任何事不做任何事offset 由 Kafka 客户端按enable.auto.commit设置自动提交ACK提交 offset提交 offset处理完成后提交REJECT_ON_ERROR提交 offset提交 offset与 ACK 相同Kafka 无原生拒信能力NACK_ON_ERROR提交 offsetseek offset回退 offset 重新消费在代码层面DefaultSubscriber与BatchSubscriber构造时会根据ack_first决定解析器行为AsyncConfluentParser(is_manualnot config.ack_first)见 faststream/confluent/subscriber/usecase.py这也是前文max_workers 1只兼容ACK_FIRST约束的底层原因之一。九、验证与测试仓库中的真实用例参考仓库中 Confluent 相关的测试集中位于 tests/brokers/confluent其中 test_config.py 覆盖了订阅者配置参数的解析与传递test_topic.py 覆盖了 Topic 对象的订阅行为。这些测试既是对本文所述参数的实证也是你编写自己的消费逻辑时可以参考的行为基线。若要本地验证基础订阅可使用 FastStream 自带的内存测试能力TestClient或启动真实 Kafka 后运行本文示例应用。小结订阅 Kafka Topic 只需broker.subscriber(topic_name)装饰一个异步函数消息按函数参数的类型注解自动解析多 Topic 可通过多个位置参数在一个订阅者内完成多个 Topic 共享一个消费者组堆叠多个装饰器则产生相互独立的处理器subscriber()背后有完整的参数体系与三种订阅者实现DefaultSubscriber/BatchSubscriber/ConcurrentDefaultSubscriberbatchTrue与max_workers1分别触发批量与并发模式使用max_workers 1时只能搭配AckPolicy.ACK_FIRST否则启动即抛SetupErrorack 策略ACK_FIRST/ACK/NACK_ON_ERROR/MANUAL决定 offset 提交时机与消费语义请根据业务对至多一次/至少一次/重试的需求选择。【免费下载链接】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),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。