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

如何向指定Kafka Topic分区发送消息?分批次指定分区实现方案

实现指定Kafka分区发送消息的两种方案

当然有可行的实现方法,针对你的需求,下面提供两种直接有效的方案:

方案一:直接在ProducerRecord中指定分区号

这是最直观的方式,利用ProducerRecord的多参数构造方法,直接指定消息要发送到的分区。修改你的代码如下:

package com.company;

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;

import java.util.Properties;

public class Main {

    public static void main(String[] args) {
        Properties properties = new Properties();
        properties.put("bootstrap.servers", "localhost:9092");
        properties.put("key.serializer", "org.apache.kafka.common.serialization.IntegerSerializer");
        properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

        Producer<Integer, String> producer = new KafkaProducer<>(properties);

        // 发送前10条到分区0,后10条到分区1(循环0-19共20条消息)
        for (int i = 0; i < 20; i++) {
            int partition = i < 10 ? 0 : 1;
            // 使用指定分区的构造方法:ProducerRecord(topic, partition, key, value)
            ProducerRecord<Integer, String> producerRecord = new ProducerRecord<>("DemoTopic", partition, i, "Test Message" + i);
            producer.send(producerRecord);
        }
        producer.close();
    }
}

说明:

  • ProducerRecord的构造方法支持显式指定分区参数,Kafka生产者会直接将消息投递到你指定的分区,绕过默认的分区规则。
  • 这里调整了循环范围为i < 20,确保正好是20条消息,前10条(0-9)发分区0,后10条(10-19)发分区1,符合你的需求。

方案二:自定义分区器(适合复杂/可复用的分区规则)

如果后续需要更复杂的分区逻辑,或者希望规则能在多个生产者中复用,可以自定义分区器。

步骤1:实现自定义分区器

package com.company;

import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.record.InvalidRecordException;

import java.util.Map;

public class CustomPartitioner implements Partitioner {

    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        // 确保key是Integer类型
        if (!(key instanceof Integer)) {
            throw new InvalidRecordException("Key must be Integer type");
        }
        Integer keyInt = (Integer) key;
        // 根据key值判断分区:小于10发分区0,否则发分区1
        return keyInt < 10 ? 0 : 1;
    }

    @Override
    public void close() {
        // 资源清理操作(如果需要)
    }

    @Override
    public void configure(Map<String, ?> configs) {
        // 配置初始化(如果需要)
    }
}

步骤2:在生产者中配置自定义分区器

修改生产者代码,添加分区器配置:

package com.company;

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;

import java.util.Properties;

public class Main {

    public static void main(String[] args) {
        Properties properties = new Properties();
        properties.put("bootstrap.servers", "localhost:9092");
        properties.put("key.serializer", "org.apache.kafka.common.serialization.IntegerSerializer");
        properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        // 指定自定义分区器
        properties.put("partitioner.class", "com.company.CustomPartitioner");

        Producer<Integer, String> producer = new KafkaProducer<>(properties);

        for (int i = 0; i < 20; i++) {
            ProducerRecord<Integer, String> producerRecord = new ProducerRecord<>("DemoTopic", i, "Test Message" + i);
            producer.send(producerRecord);
        }
        producer.close();
    }
}

说明:

  • 自定义分区器实现了Partitioner接口,重写partition方法定义自己的分区逻辑。
  • 这种方式适合需要统一管理分区规则的场景,多个生产者可以复用同一个分区器。

总结

  • 如果只是简单的固定规则分区,方案一更直接高效,无需额外代码。
  • 如果需要复杂的、可复用的分区逻辑,方案二更合适。

内容的提问来源于stack exchange,提问作者Abhishek

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 19:21:35