如何在Java代码中实现Kafka同一消费组内的多消费者
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.strategyconfig toorg.apache.kafka.clients.consumer.RoundRobinAssignor. - Graceful Shutdown: Always clean up consumers properly (call
wakeup()andclose()) to avoid leaving connections open or losing offset commits.
内容的提问来源于stack exchange,提问作者B.Y.

