Helidon SE 4.1.6中使用Kafka Producer发送数据至指定分区
在Helidon SE 4.1.6中实现Kafka指定分区消息发送
问题描述
我在Helidon SE 4.1.6环境下,需要通过Kafka Producer将数据发送到Apache Kafka的指定分区。目前已实现基础的主题消息发送,但尝试通过ProducerRecord指定分区发送时,数据并未送达目标分区,通过消费者验证确认消息未出现在预期分区中。
可正常工作的基础发送代码
public class KafkaProducer { public static void init() { String kafkaServer = "localhost:9092"; String topic = "topic2"; Channel<String> toKafka = Channel.<String>builder() .subscriberConfig(KafkaConnector.configBuilder() .bootstrapServers(kafkaServer) .topic(topic) .keySerializer(StringSerializer.class) .valueSerializer(StringSerializer.class) .build()) .build(); KafkaConnector kafkaConnector = KafkaConnector.create(); System.out.println("In producer: "); Messaging messaging = Messaging.builder() .publisher(toKafka, Multi.just("Test1", "Test2").map(Message::of)) .connector(kafkaConnector) .build() .start(); } }
尝试的指定分区代码(未生效)
public class KafkaProducerPar { public static void init() { String bootstrapServers = "localhost:9092"; String topic = "topic3"; // Create a KafkaConnector Config config = Config.create(); KafkaConnector kafkaConnector = KafkaConnector.create(config); // Create a Channel for producing messages Channel<ProducerRecord<String, String>> toKafka = Channel.<ProducerRecord<String, String>>builder() .subscriberConfig(KafkaConnector.configBuilder() .bootstrapServers(bootstrapServers) .topic(topic) .keySerializer(StringSerializer.class) .valueSerializer(StringSerializer.class) .build()) .build(); // Create a Messaging instance Messaging messaging = Messaging.builder() .publisher(toKafka,createMessageStream1(topic) ) .connector(kafkaConnector) .build() .start(); Messaging messaging2 = Messaging.builder() .publisher(toKafka, Multi.just( Message.of(new ProducerRecord<>(topic, 0, "key1", "Message for partition 0")), Message.of(new ProducerRecord<>(topic, 1, "key2", "Message for partition 1")) )) .connector(kafkaConnector) .build() .start(); } }
问题原因及修复方案
核心问题
Helidon Reactive Messaging的Kafka连接器中,若在Channel配置里指定了topic参数,会覆盖ProducerRecord中自定义的主题、分区等信息,导致你的分区设置完全失效。此外原代码未等待消息发送完成就可能退出程序,也会造成消息丢失。
修复步骤
- 移除Channel的topic配置:让
ProducerRecord自行指定目标主题和分区 - 等待消息发送完成:添加等待逻辑确保消息被Kafka接收
- 简化Messaging实例创建:无需重复创建多个Messaging实例
修复后的完整代码
import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import io.helidon.messaging.Channel; import io.helidon.messaging.Messaging; import io.helidon.messaging.kafka.KafkaConnector; import io.smallrye.mutiny.Multi; import io.helidon.messaging.Message; public class KafkaProducerPar { public static void init() throws InterruptedException { String bootstrapServers = "localhost:9092"; String topic = "topic3"; // 创建Kafka连接器 KafkaConnector kafkaConnector = KafkaConnector.create(); // 关键:Channel不指定topic,由ProducerRecord决定目标 Channel<ProducerRecord<String, String>> toKafka = Channel.<ProducerRecord<String, String>>builder() .subscriberConfig(KafkaConnector.configBuilder() .bootstrapServers(bootstrapServers) .keySerializer(StringSerializer.class) .valueSerializer(StringSerializer.class) .build()) .build(); // 构造指定分区的消息流 Multi<Message<ProducerRecord<String, String>>> messageStream = Multi.just( Message.of(new ProducerRecord<>(topic, 0, "key1", "Message for partition 0")), Message.of(new ProducerRecord<>(topic, 1, "key2", "Message for partition 1")) ); // 启动Messaging并等待发送完成 Messaging messaging = Messaging.builder() .publisher(toKafka, messageStream) .connector(kafkaConnector) .build() .start(); // 等待消息发送完成(根据业务场景调整等待时长) messaging.awaitTermination(java.time.Duration.ofSeconds(5)); } public static void main(String[] args) throws InterruptedException { init(); } }
验证建议
- 确认目标主题
topic3的分区数≥2(你指定了分区0和1) - 使用Kafka命令行工具直接验证消息分区:
该命令会打印每条消息的所属分区,可直观确认是否发送到指定分区。kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic topic3 --property print.partition=true --from-beginning
内容的提问来源于stack exchange,提问作者MOHAMMAD SHADAB
相关产品推荐
相关产品推荐

