ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

FastStream 中通过分区键(Partition Key)控制 Kafka 消息路由:发布、读取与底层实现

FastStream 中通过分区键(Partition Key)控制 Kafka 消息路由:发布、读取与底层实现 FastStream 中通过分区键Partition Key控制 Kafka 消息路由发布、读取与底层实现【免费下载链接】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导读本指南围绕 FastStream 官方文档中 Using a Partition Key 一章展开讲解如何在使用KafkaBroker.publisher(...)装饰器发布消息时通过key参数为消息指定 Kafka 分区键从而控制同一 Key 的消息进入同一分区实现有序处理与负载均衡。读完本文你将掌握在 FastStream 中定义带 Key 的发布者、在publish调用中传递key、在订阅端读取消息 Key以及 Key 的序列化与分区器Partitioner的底层工作机制。为什么分区键Partition Key如此重要在 ApacheKafka中分区键是一个核心概念它决定了消息最终被写入哪个分区Partition。借助分区键可以实现两个关键目标保持相关消息的局部有序把具有相同 Key 的消息发送到同一个分区Kafka 仅在分区内部保证消息顺序因此同 Key 消息聚集在同一分区天然维持处理顺序适合订单事件、用户会话等强顺序场景横向扩展与负载均衡Kafka 通过分区将消息分散到多个 Broker 上并行处理同时数据跨 Broker 副本复制提供了容错能力。FastStream 允许你在KafkaBroker.publisher(...)装饰器创建的发布者上指定分区键。下文将完整演示具体做法。发布带分区键的消息两步走原文档给出的流程分为两步下面结合 docs/docs_src/kafka/publish_with_partition_key/app.py 中的真实示例代码逐一展开。Step 1定义 Publisher在你的 FastStream 应用中通过KafkaBroker.publisher(...)装饰器定义发布者。该装饰器允许配置发布相关的各种属性其中就包括分区键相关的参数。下面的代码先创建了一个指向output_data主题的发布者对象from pydantic import BaseModel, Field, NonNegativeFloat from faststream import Context, FastStream, Logger from faststream.kafka import KafkaBroker class Data(BaseModel): data: NonNegativeFloat Field( ..., examples[0.5], descriptionFloat data example, ) broker KafkaBroker(localhost:9092) app FastStream(broker) to_output_data broker.publisher(output_data)关键点说明broker KafkaBroker(localhost:9092)创建指向本地 Kafka 实例的 Broker地址可替换为你的实际 Broker 地址app FastStream(broker)将 Broker 挂载到 FastStream 应用to_output_data broker.publisher(output_data)返回一个可复用的发布者对象后续可以多次调用它的publish方法。Step 2调用 publish 时传入 Key当准备好向指定主题发布消息时只需要在publish调用中增加key参数即可。该参数会被用来确定消息应进入哪个分区await to_output_data.publish(Data(datamsg.data 1.0), keybkey)key的类型可以是bytes、任意可被key_serializer序列化为字节的对象或None。传入None时由分区器随机选择分区详见下文底层原理一节。完整示例应用消费后带 Key 转发下面是一段完整应用它从input_data主题消费消息处理后携带指定 Key 发布到output_data主题并演示如何在订阅端读取到消息的 Key。from pydantic import BaseModel, Field, NonNegativeFloat from faststream import Context, FastStream, Logger from faststream.kafka import KafkaBroker class Data(BaseModel): data: NonNegativeFloat Field( ..., examples[0.5], descriptionFloat data example, ) broker KafkaBroker(localhost:9092) app FastStream(broker) to_output_data broker.publisher(output_data) broker.subscriber(input_data) async def on_input_data( msg: Data, logger: Logger, key: bytes Context(message.raw_message.key), ) - None: logger.info(on_input_data(msg%s), msg) await to_output_data.publish(Data(datamsg.data 1.0), keybkey) broker.subscriber(output_data) async def on_output_data( msg: Data, logger: Logger, key: bytes Context(message.raw_message.key), ) - None: logger.info(on_output_data(msg%s), msg)与标准发布相比唯一的区别就是在publish调用中增加了key参数——这正是控制 Kafka 分区与消息处理方式的关键所在。代码中的额外细节key: bytes Context(message.raw_message.key)通过 FastStream 的Context注入从原始 Kafka 消息中取出 Key 供订阅处理器使用on_output_data同样读取了 Key说明带 Key 发布的消息在被消费时依然携带 Key 信息Data(datamsg.data 1.0)展示了在转发时对消息体进行加工的场景。深入源码key 参数的完整语义publish 调用链中的 key从 faststream/kafka/publisher/usecase.py 的DefaultPublisher.publish实现可以看到key是发布方法的标准关键字参数async def publish( self, message: SendableMessage, topic: str , *, key: bytes | Any | None None, partition: int | None None, ... ) - Union[asyncio.Future[RecordMetadata], RecordMetadata]:其文档字符串对key的语义描述为用于关联消息的键可用来决定消息发往哪个分区当partition为None且生产者分区器保持默认配置时相同 Key 的消息会被投递到同一分区若key为None则分区随机选择。Key 必须是bytes类型或可以通过配置的key_serializer序列化为字节。值得注意的两个实现细节发布者默认 KeyKafkaPublisherConfig见 faststream/kafka/publisher/config.py中包含key: bytes | str | None字段DefaultPublisher初始化时将其保存为self.key调用级 Key 优先级cmd KafkaPublishCommand(..., keykey or self.key, ...)表明publish(key...)中显式传入的 Key 会覆盖发布者配置的默认 Key_publish路径中同样存在cmd.key cmd.key or self.key的回退逻辑。也就是说除了在每次publish调用时传 Key你还可以在创建发布者时配置一个默认 Key未显式传入时自动生效。分区器与 Key 的哈希规则在 faststream/kafka/broker/broker.py 的 Broker 初始化参数中可以看到两个与 Key 直接相关的底层配置key_serializer用于将用户提供的 Key 转换为bytespartitioner决定每条消息被分配到哪个分区的可调用对象。默认分区器使用与 Java 客户端一致的murmur2 哈希算法对每个非NoneKey 取哈希从而保证相同 Key 的消息被分配到同一分区当 Key 为None时消息被投递到随机分区。此外publish还支持partition参数当需要完全绕过 Key 路由、直接指定分区时可以显式传入partition此时分区选择不再依赖 Key 的哈希结果调用级partition优先于发布者配置的partition。请求-响应RPC与批量发布中的 KeyLogicPublisher.request(...)与DefaultPublisher.request(...)同样接受key参数因此基于 Kafka 的请求-响应RPC模式也能携带分区键并同样遵循key or self.key的默认回退规则BatchPublisher.publish(*messages, ...)支持为一批消息指定同一个 Key当未显式指定partition且使用默认分区器时整批消息会被路由到同一分区从而保证批内消息的连续性。订阅端如何读取 KeyKafka 消息的 Key 在消费侧同样有意义。FastStream 提供了两种途径Context 注入如示例所示使用key: bytes Context(message.raw_message.key)直接取原始消息中的 Key这是最常见的用法key_deserializer在 faststream/kafka/broker/registrator.py 中订阅注册函数如subscriber、batch_subscriber都提供key_deserializer: Callable[[bytes], Any] | None None参数用于将字节形式的 Key 反序列化为业务类型例如解析为字符串或自定义对象与value_deserializer对消息体的处理相对应。用 TestKafkaBroker 验证带 Key 的发布仓库为本文示例提供了配套测试 tests/docs/kafka/publish_with_partition_key/test_app.py展示了如何在不依赖真实 Kafka 集群的情况下验证带 Key 的发布流程async with TestKafkaBroker(broker): await broker.publish(Data(data0.2), input_data, keybmy_key) on_input_data.mock.assert_called_once_with(dict(Data(data0.2))) to_output_data.mock.assert_called_once_with(dict(Data(data1.2)))通过TestKafkaBroker(broker)上下文管理器替换真实 Brokerbroker.publish(..., keybmy_key)模拟外部生产者向input_data发布带 Key 的消息on_input_data.mock与to_output_data.mock分别断言消费与转发逻辑按预期触发0.2被加1.0后变成1.2测试文件中还保留了一个被skip的test_keys用例注释明确说明应能带 Key 发布消息且应校验 Key 本身——它展示了更严格地断言key被正确透传的测试思路可以作为你编写自己的 Key 断言测试的参考。总结在 Kafka 中使用分区键是优化消息分布、维持消息顺序、实现高效处理的基础实践。在 FastStream 中这一能力被浓缩为两个动作定义发布者、在publish调用中传入key。底层默认采用 murmur2 哈希将相同 Key 路由到同一分区key_serializer负责 Key 的字节化partitioner与partition参数则提供了更细粒度的分区控制订阅端既可以通过Context直接读取 Key也可以通过key_deserializer将其反序列化。配合TestKafkaBroker的 mock 断言你可以安全、快速地为基于 Kafka 的 FastStream 应用加上分区键能力让消息处理按业务维度有序、均衡地横向扩展。【免费下载链接】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),仅供参考
返回列表