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

如何在Java代码中实现Kafka同一消费组内的多消费者

Setting Up Multiple Kafka Consumers in the Same Group (Java)

Got it, let's tackle this—since you already know the CLI approach works, translating that to Java code is straightforward once you remember the core Kafka consumer group rules. Here's how to set up 2 consumers in the same group for your 6-partition topic:

Core Background

The key rule here is: all consumers in the same group must use the exact same group.id. Kafka will automatically distribute the topic's partitions across active consumers in the group. With 6 partitions and 2 consumers, you'll end up with 3 partitions assigned to each consumer (using the default RangeAssignor strategy).


Option 1: Using the Native Kafka Client

If you're working directly with the Kafka Java client, you'll need to spin up separate threads for each consumer (each with the same group ID) and handle their lifecycle.

Step 1: Define Consumer Configurations

First, create a shared properties object with the critical group.id set:

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.util.Properties;

private static Properties getSharedConsumerProps() {
    Properties props = new Properties();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-broker:9092");
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "your-unique-group-id"); // Same for all consumers!
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
    // Optional: Disable auto-commit if you want manual offset control
    // props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
    return props;
}

Step 2: Create a Reusable Consumer Runnable

Wrap the consumer logic in a Runnable so you can spawn multiple threads:

import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.time.Duration;
import java.util.Collections;

public class KafkaConsumerWorker implements Runnable {
    private final KafkaConsumer<String, String> consumer;
    private final String topic;
    private volatile boolean isRunning = true;

    public KafkaConsumerWorker(Properties props, String topic) {
        this.consumer = new KafkaConsumer<>(props);
        this.topic = topic;
    }

    @Override
    public void run() {
        consumer.subscribe(Collections.singletonList(topic));
        while (isRunning) {
            ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
            // Process your messages here
            records.forEach(record -> {
                System.out.printf("Thread: %s | Partition: %d | Value: %s%n",
                        Thread.currentThread().getName(), record.partition(), record.value());
            });
            // If using manual commit:
            // consumer.commitSync();
        }
        consumer.close();
    }

    // Call this for graceful shutdown
    public void stop() {
        isRunning = false;
        consumer.wakeup();
    }
}

Step 3: Launch Multiple Consumers

Spawn two threads, each running an instance of your consumer worker:

public class MultiConsumerLauncher {
    public static void main(String[] args) {
        String targetTopic = "your-6-partition-topic";
        Properties sharedProps = getSharedConsumerProps();

        // Create two consumer workers
        KafkaConsumerWorker consumer1 = new KafkaConsumerWorker(sharedProps, targetTopic);
        KafkaConsumerWorker consumer2 = new KafkaConsumerWorker(sharedProps, targetTopic);

        // Start them in separate threads
        Thread thread1 = new Thread(consumer1, "Consumer-1");
        Thread thread2 = new Thread(consumer2, "Consumer-2");
        thread1.start();
        thread2.start();

        // Add a shutdown hook to clean up properly
        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            consumer1.stop();
            consumer2.stop();
            try {
                thread1.join();
                thread2.join();
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }));
    }
}

Option 2: Using Spring Kafka (Simpler for Spring Apps)

If you're using Spring Boot with Spring Kafka, it's even easier—you just need to configure the concurrency level, and Spring handles spawning the consumer threads for you.

Step 1: Configure Application Properties

In application.yml (or application.properties), set the group ID and concurrency:

spring:
  kafka:
    bootstrap-servers: your-kafka-broker:9092
    consumer:
      group-id: your-unique-group-id # Same group ID here
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      auto-offset-reset: latest
    listener:
      concurrency: 2 # Number of consumer threads to spawn

Step 2: Create a Listener Method

Just define a single listener method—Spring will run it across 2 threads, each handling its assigned partitions:

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;

@Service
public class TopicConsumerService {

    @KafkaListener(topics = "your-6-partition-topic")
    public void consumeMessage(ConsumerRecord<String, String> record) {
        // Your message processing logic here
        System.out.printf("Thread: %s | Partition: %d | Value: %s%n",
                Thread.currentThread().getName(), record.partition(), record.value());
    }
}

Key Notes

  • Group ID Consistency: Double-check that all consumers in the group use the exact same group.id—this is non-negotiable for partition assignment.
  • Partition vs Consumer Count: Since you have 6 partitions and 2 consumers, each will get 3 partitions. If you had more consumers than partitions, the extra consumers would sit idle until a running consumer drops out.
  • Partition Assignment Strategy: If you want to change how partitions are split (e.g., round-robin instead of range), set the partition.assignment.strategy config to org.apache.kafka.clients.consumer.RoundRobinAssignor.
  • Graceful Shutdown: Always clean up consumers properly (call wakeup() and close()) to avoid leaving connections open or losing offset commits.

内容的提问来源于stack exchange,提问作者B.Y.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.27 19:32:29