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

未指定分区时Kafka消费者无法消费消息,如何在Java KafkaConsumer中指定分区?

Fixing Kafka Consumer: How to Specify a Partition in Java

Hey there! I see you're running into an issue where your Kafka consumer only works when explicitly targeting partition 0 via the CLI, and your Java consumer using subscribe() isn't picking up any records. Let's break down how to fix this by manually specifying the partition in your Java code, plus a few checks to understand why this is happening.

Why This Happens

When you use kafka-console-consumer without --partition, Kafka's consumer coordinator tries to assign partitions to your consumer group automatically. If your consumer group my-kafka-java-app has already committed an offset that's at the end of the partition (or if the auto-assignment logic isn't picking up partition 0 for some reason), you won't see any records. Using --partition 0 skips this auto-assignment and directly reads from that specific partition, which is why it works.

How to Specify a Partition in Java KafkaConsumer

Instead of using subscribe() (which relies on group-based partition assignment), you can use assign() to manually target the specific partition you want to consume from. Here's how to modify your code:

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class ConsumerDemo {
    public static void main(String[] args) {
        Logger logger = LoggerFactory.getLogger(ConsumerDemo.class.getName());
        String bootstrapServer = "127.0.0.10:9092";
        String groupId = "my-kafka-java-app";
        String topic = "first-topic";
        int targetPartition = 0; // The partition you want to consume

        // create consumer configs
        Properties properties = new Properties();
        properties.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer);
        properties.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        properties.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        properties.setProperty(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        properties.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        // create consumer
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(properties);

        // Manually assign the target partition instead of subscribing to the topic
        TopicPartition partition = new TopicPartition(topic, targetPartition);
        consumer.assign(Collections.singletonList(partition));

        // poll for new data
        while(true){
            ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
            for (ConsumerRecord<String, String> record : records){
                logger.info("Key: " + record.key() + ", Value: "+ record.value() );
                logger.info("Partition: " + record.partition() + ", Offset: "+ record.offset());
            }
        }
    }
}

Key Differences Between subscribe() and assign()

  • subscribe(): Works with consumer groups, letting Kafka automatically assign partitions across group members. This is great for scaling consumers across multiple partitions.
  • assign(): Lets you manually specify exactly which partitions to consume. This bypasses the consumer group coordinator, so you won't get automatic partition rebalancing if other consumers join/leave the group. Use this when you need direct control over which partition your consumer reads from.

Additional Checks to Debug

  1. Check your consumer group's offset position: Run this CLI command to see where your group's offset is for the topic:

    kafka-consumer-groups.sh --bootstrap-server 127.0.0.10:9092 --describe --group my-kafka-java-app
    

    If the CURRENT-OFFSET matches the LOG-END-OFFSET for partition 0, that means your group has already consumed all available records. Even with auto.offset.reset=earliest, Kafka won't re-read old records unless you reset the offset.

  2. Verify topic partition count: Make sure your topic actually has a partition 0 (it should by default, but it's worth checking):

    kafka-topics.sh --bootstrap-server 127.0.0.10:9092 --describe --topic first-topic
    

With these changes, your Java consumer should now read directly from partition 0, just like your working CLI command does.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.01 00:47:38