在 Ubuntu 上配置 Kafka 消费组(Consumer Group),通常分几种场景:命令行消费、Java 客户端配置、以及 Kafka 配置文件层面的设置。下面按常见用法说明。
假设你已经安装好 Kafka(如 /opt/kafka 或 /usr/local/kafka),并且:
kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic test-topic \
--group my-consumer-group \
--from-beginning
说明:
--group:消费组名称(不存在会自动创建)kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group my-consumer-group
输出示例:
TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
test 0 10 15 5
kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--group my-consumer-group \
--reset-offsets \
--to-earliest \
--execute \
--topic test-topic
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.6.0</version>
</dependency>
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "my-consumer-group"); // ✅ 消费组
props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "5000");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("test-topic"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.println(record.value());
}
}
✅ group.id 就是消费组
一般不需要特别配置消费组,但常见相关参数:
# 消费组 offset 存储(老版本)
offsets.topic.replication.factor=1
# 新版本默认使用内部 topic __consumer_offsets
✅ 同一 group.id 的消费者
✅ 不同 group.id
✅ 消费组卡住
--groupgroup.idcd /opt/kafka/bin
或加到 PATH:
export PATH=$PATH:/opt/kafka/bin
如果你需要:
告诉我你的使用场景,我可以直接给你完整配置。