ARTICLE DETAIL

资讯详情

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

FastStream Confluent 批量发布(Batch Publishing)实战:从 `batch=True` 到 `KafkaPublishMessage` 逐条属性控制

FastStream Confluent 批量发布(Batch Publishing)实战:从 `batch=True` 到 `KafkaPublishMessage` 逐条属性控制 FastStream Confluent 批量发布Batch Publishing实战从batchTrue到KafkaPublishMessage逐条属性控制【免费下载链接】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 的faststream.confluent模块深入讲解如何使用broker.publisher(..., batchTrue)装饰器将多条消息一次性批量发送到 Confluent Kafka 主题覆盖「创建批量 Publisher」「返回元组自动触发批量发送」「直接调用publish_batch(...)」三种核心用法并重点剖析KafkaPublishMessage如何为同一批消息中的每一条设置独立的 key、headers、timestamp 与 correlation_id。读完本文你将能够写出吞吐更高、网络开销更低、且支持按分区精确路由的 Kafka 批量生产代码。一、General Overview批量发布的整体思路在 FastStream 的 Confluent 集成中如果你的业务需要把多条数据一次性发送出去#!python broker.publisher(...)装饰器就提供了便捷的批量能力。要启用批量生产只需完成两个关键步骤创建 Publisher 时设置batchTrue这一配置告诉 Publisher你打算以批量模式发送消息在生产者函数中返回一个由消息组成的元组tuple该返回动作会触发生产者把元组中的消息收集起来作为一批一次性发送给Kafkabroker。下面以一个详细示例说明如何一边从input_data_1主题消费一边把处理结果批量生产到output_data主题。补充说明batchTrue不仅改变了 Publisher 的「语义」在底层它还会被实例化为独立的BatchPublisher类见 faststream/confluent/publisher/usecase.py其publish方法签名从「单条消息」变为「可变参数*messages」并统一走_basic_publish_batch的批量发送链路。二、完整代码示例从订阅到批量生产的应用先看完整的应用创建过程随后再拆解批量生产的各个步骤。以下是示例应用的完整代码源码位于 docs/docs_src/confluent/publish_batch/app.pyfrom typing import Tuple from pydantic import BaseModel, Field, NonNegativeFloat from faststream import FastStream, Logger from faststream.confluent import KafkaBroker class Data(BaseModel): data: NonNegativeFloat Field( ..., examples[0.5], descriptionFloat data example, ) broker KafkaBroker(localhost:9092) app FastStream(broker) decrease_and_increase broker.publisher(output_data, batchTrue) decrease_and_increase broker.subscriber(input_data_1) async def on_input_data_1(msg: Data, logger: Logger) - Tuple[Data, Data]: logger.info(msg) return Data(data(msg.data * 0.5)), Data(data(msg.data * 2.0)) broker.subscriber(input_data_2) async def on_input_data_2(msg: Data, logger: Logger) - None: logger.info(msg) await decrease_and_increase.publish( Data(data(msg.data * 0.5)), Data(data(msg.data * 2.0)), )这个应用同时实现了批量发布的两种触发方式触发方式关键代码适用场景装饰器 返回元组decrease_and_increase修饰消费函数函数return两个Data处理函数本身天然产出多条结果无需手动调用直接调用.publish(...)在on_input_data_2中await decrease_and_increase.publish(Data(...), Data(...))需要在中途主动推送批量数据的任意代码路径两种方式最终都会把两条Data记录作为一批发送到output_data主题。Step 1创建批量 Publisherdecrease_and_increase broker.publisher(output_data, batchTrue)这一行声明了一个名为decrease_and_increase的批量 Publisher目标主题是output_databatchTrue是关键配置。后续它可以同时扮演「装饰器」和「可直接调用的对象」两种角色。Step 2发布真正的一批消息方式 A直接调用 Publisher 发布批量消息await decrease_and_increase.publish( Data(data(msg.data * 0.5)), Data(data(msg.data * 2.0)), )方式 B装饰处理函数并返回批量消息decrease_and_increase broker.subscriber(input_data_1) async def on_input_data_1(msg: Data, logger: Logger) - Tuple[Data, Data]: return Data(data(msg.data * 0.5)), Data(data(msg.data * 2.0))示例应用把这两种方式都实现了你可以按需选择更适合自己业务形态的一种。直接通过 broker 对象批量发布除了 Publisher 对象你还可以直接从broker上调用publish_batch方法。例如#!python broker.publish_batch(msg1, msg2, topicoutput_data)即可不创建任何中间 Publisher直接把多条消息发送到指定主题对应源码为 faststream/confluent/broker/broker.py 中的publish_batch方法。测试验证批量发布的行为是确定的该示例的批量行为有对应的单元测试见 tests/docs/confluent/publish_batch/test_app.py使用TestKafkaBroker进行内存级验证async with TestKafkaBroker(broker): await broker.publish(Data(data2.0), input_data_1) on_input_data_1.mock.assert_called_once_with(dict(Data(data2.0))) decrease_and_increase.mock.assert_called_once_with( [dict(Data(data1.0)), dict(Data(data4.0))], )测试断言表明消费到Data(data2.0)后批量 Publisher 收到的是[1.0, 4.0]这样一条批量记录列表形式而不是两次独立的单条发布。这也从测试层面印证了batchTrue的语义。三、批量发布内部的调用链消息如何被组装成一批为了真正理解批量发布可以顺着源码看一次publish_batch的完整链路命令构造KafkaPublishCommand接受*messages可变参数把多条消息统一放进batch_bodies并解析出逐条 key见 faststream/confluent/response.py批量发送入口无论来自 Publisher 对象还是 broker 对象最终都会调用_basic_publish_batch由AsyncConfluentFastProducerImpl.publish_batch执行见 faststream/confluent/publisher/producer.py编码底层先通过codec默认DefaultCodec逐条编码batch_bodies若传入的 codec 实现了BatchCodecProto则会调用encode_batch一次性编码整批数据组装批次调用 confluent-kafka 的producer.create_batch()创建批次对象逐条batch.append(...)附上各自的消息体、key、timestamp 与 headers发送最后通过producer.send_batch(batch, destination, partition..., no_confirm...)一次性发给 broker。其中第 4 步用到的逐条 key 来自cmd.key_for(index)其实现会优先使用该条消息自身的 key否则回退到publish_batch调用级传入的默认 key实现见 faststream/_internal/kafka/keys.py 中的key_for_index。这正是下一节「逐条属性控制」的底层基础。四、Per-Message Attributes用KafkaPublishMessage为每条消息设置独立属性批量发布时你常常需要为同一批里的每一条消息分配不同的 key、headers 或 timestamp。为了支持这一场景FastStream 提供了#!python KafkaPublishMessage辅助类型——它是#!python KafkaResponse的语义别名见 faststream/confluent/response.py 末尾的KafkaPublishMessage KafkaResponse专门用于在#!python publish_batch(...)调用内部构造带属性的消息。通过把 payload 包装进#!python KafkaPublishMessage你可以为单条消息附加以下属性属性类型说明keybytes \| Any \| NoneKafka 消息 key用于分区路由headersdict[str, Any] \| None单条消息的自定义 headerstimestamp_msint \| None显式指定的消息时间戳correlation_idstr \| None自定义关联标识自由混用带属性的消息与普通消息共存你可以在同一次调用中自由混用#!python KafkaPublishMessage实例和原始 payload。纯值字符串、bytes、dict、模型等会原样发布并使用从#!python publish_batch(...)调用继承来的默认 key默认是#!python Nonefrom faststream.confluent import KafkaBroker, KafkaPublishMessage broker KafkaBroker() broker.subscriber(input) async def handler() - None: await broker.publish_batch( KafkaPublishMessage(user:1, keybuser1), KafkaPublishMessage(user:2, keybuser2), user:3, # Uses default key (None) topicoutput, )在上面的例子中前两条消息携带各自的专用 keybuser1和buser2第三条以纯字符串传入的消息则回退到默认 key。底层的 key 对齐机制从源码看这一混用能力由 faststream/_internal/kafka/keys.py 支撑extract_per_message_keys_and_bodies通过singledispatch识别Response对象KafkaResponse是Response的子类从每条消息中抽取 body 与 key只有存在至少一个非None的逐条 key 时才会做归一化处理否则直接复用原始批量数据、避免额外分配。随后KafkaPublishCommand把这些逐条 key 与batch_bodies保持对齐即使批量内容在后续流程中被改写如空批量填充realign_keys也会同步修正 key 的对应关系。命名说明KafkaPublishMessage是在publish_batch(...)中构造出站消息时推荐的、更语义化的名字底层它与KafkaResponse是同一个对象两个名字完全可互换。KafkaResponse除了用于发布外还可以作为 handler 的返回值直接发送单条响应消息。实践建议当同一批消息需要经由不同的 key 路由到不同分区或需要为单条消息附加元数据时优先使用KafkaPublishMessage——这是在不把批量拆成多次独立publish(...)调用的前提下控制单条消息属性的最干净方式。五、为什么要批量发布通过上面的示例你已经掌握了如何借助broker.publisher(..., batchTrue)在 FastStream 与 Confluent Kafka 中高效地批量发布消息。遵循前文提到的两个关键步骤可以显著提升基于 Kafka 的应用的性能与可靠性。具体而言批量发布在 Kafka 场景下有如下优势更高的吞吐Improved Throughput批量发布允许在一次传输中发送多条消息减少了逐条投递带来的开销从而提升 Kafka 应用的吞吐量并降低延迟。降低网络与 broker 负载Reduced Network and Broker Load批量发送减少了网络调用次数与 broker 交互次数让 Kafka 集群与网络资源都更高效。原子性Atomicity批量机制保证一组相关消息要么一起被处理、要么都不处理。在需要维持数据一致性与完整性的处理场景中这种原子性至关重要。更强的可扩展性Enhanced Scalability借助批量发布你可以更高效地扩展 Kafka 应用以支撑高消息量。以更大的块发送消息能更充分地利用 Kafka 的并行度与分区能力。一点补充说明在底层实现中批量发布一次只产生一个Kafka 批次请求见 faststream/confluent/publisher/producer.py 的create_batch/send_batch这正是吞吐提升的直接来源同时批量的原子性也体现在批量请求级别的统一发送语义上。不过需要留意的是跨主题、跨分区的「事务性原子性」属于 Kafka 事务transactions能力的范畴与本文的批量发布是两个不同概念请勿混淆。六、总结在 FastStream 的 Confluent 集成中批量发布是一个「两步行」的能力创建 Publisher 时设置batchTrue返回元组或直接调用publish(...)/broker.publish_batch(...)传入多条消息。当需要为同一批内的每条消息设置独立 key、headers、timestamp 或 correlation_id 时使用KafkaPublishMessage包装对应消息即可它既可以与普通 payload 自由混用又能在不拆分批量调用的前提下实现精确的逐条控制。相关完整示例与测试分别位于 docs/docs_src/confluent/publish_batch/app.py 与 tests/docs/confluent/publish_batch/test_app.py你可以直接复制示例、配合TestKafkaBroker在本地无 broker 环境中验证整条批量发布链路。【免费下载链接】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),仅供参考
返回列表