You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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中自定义的主题、分区等信息,导致你的分区设置完全失效。此外原代码未等待消息发送完成就可能退出程序,也会造成消息丢失。

修复步骤

  1. 移除Channel的topic配置:让ProducerRecord自行指定目标主题和分区
  2. 等待消息发送完成:添加等待逻辑确保消息被Kafka接收
  3. 简化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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 19:55:56