FastStream Confluent 批量发布(Batch Publishing)实战:从batch=True到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(..., batch=True)装饰器将多条消息一次性批量发送到 Confluent Kafka 主题,覆盖「创建批量 Publisher」「返回元组自动触发批量发送」「直接调用publish_batch(...)」三种核心用法,并重点剖析KafkaPublishMessage如何为同一批消息中的每一条设置独立的 key、headers、timestamp 与 correlation_id。读完本文,你将能够写出吞吐更高、网络开销更低、且支持按分区精确路由的 Kafka 批量生产代码。
一、General Overview:批量发布的整体思路
在 FastStream 的 Confluent 集成中,如果你的业务需要把多条数据一次性发送出去,#!python @broker.publisher(...)装饰器就提供了便捷的批量能力。要启用批量生产,只需完成两个关键步骤:
- 创建 Publisher 时设置
batch=True:这一配置告诉 Publisher,你打算以批量模式发送消息; - 在生产者函数中返回一个由消息组成的元组(tuple):该返回动作会触发生产者把元组中的消息收集起来,作为一批一次性发送给Kafkabroker。
下面以一个详细示例说明:如何一边从"input_data_1"主题消费,一边把处理结果批量生产到"output_data"主题。
补充说明:
batch=True不仅改变了 Publisher 的「语义」,在底层它还会被实例化为独立的BatchPublisher类(见 faststream/confluent/publisher/usecase.py),其publish方法签名从「单条消息」变为「可变参数*messages」,并统一走_basic_publish_batch的批量发送链路。
二、完整代码示例:从订阅到批量生产的应用
先看完整的应用创建过程,随后再拆解批量生产的各个步骤。以下是示例应用的完整代码(源码位于 docs/docs_src/confluent/publish_batch/app.py):
from 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], description="Float data example", ) broker = KafkaBroker("localhost:9092") app = FastStream(broker) decrease_and_increase = broker.publisher("output_data", batch=True) @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:创建批量 Publisher
decrease_and_increase = broker.publisher("output_data", batch=True)这一行声明了一个名为decrease_and_increase的批量 Publisher,目标主题是"output_data",batch=True是关键配置。后续它可以同时扮演「装饰器」和「可直接调用的对象」两种角色。
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", topic="output_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(data=2.0), "input_data_1") on_input_data_1.mock.assert_called_once_with(dict(Data(data=2.0))) decrease_and_increase.mock.assert_called_once_with( [dict(Data(data=1.0)), dict(Data(data=4.0))], )测试断言表明:消费到Data(data=2.0)后,批量 Publisher 收到的是[1.0, 4.0]这样一条批量记录(列表形式),而不是两次独立的单条发布。这也从测试层面印证了batch=True的语义。
三、批量发布内部的调用链:消息如何被组装成一批
为了真正理解批量发布,可以顺着源码看一次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,你可以为单条消息附加以下属性:
| 属性 | 类型 | 说明 |
|---|---|---|
key | bytes \| Any \| None | Kafka 消息 key,用于分区路由 |
headers | dict[str, Any] \| None | 单条消息的自定义 headers |
timestamp_ms | int \| None | 显式指定的消息时间戳 |
correlation_id | str \| None | 自定义关联标识 |
自由混用:带属性的消息与普通消息共存
你可以在同一次调用中自由混用#!python KafkaPublishMessage实例和原始 payload。纯值(字符串、bytes、dict、模型等)会原样发布,并使用从#!python publish_batch(...)调用继承来的默认 key(默认是#!python None):
from faststream.confluent import KafkaBroker, KafkaPublishMessage broker = KafkaBroker() @broker.subscriber("input") async def handler() -> None: await broker.publish_batch( KafkaPublishMessage("user:1", key=b"user1"), KafkaPublishMessage("user:2", key=b"user2"), "user:3", # Uses default key (None) topic="output", )在上面的例子中,前两条消息携带各自的专用 key(b"user1"和b"user2"),第三条以纯字符串传入的消息则回退到默认 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(..., batch=True)在 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 时设置
batch=True; - 返回元组,或直接调用
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),仅供参考