下面以 Debian 系统为背景,给你一个 Kafka 消费者配置与使用示例,分为 命令行方式 和 代码方式(Java/Python),看你用什么场景。
sudo apt update
sudo apt install default-jdk -y
下载 Kafka(示例版本 3.6.1):
wget https://downloads.apache.org/kafka/3.6.1/kafka_2.13-3.6.1.tgz
tar -xzf kafka_2.13-3.6.1.tgz
cd kafka_2.13-3.6.1
# 启动 Zookeeper(老版本)
bin/zookeeper-server-start.sh config/zookeeper.properties
# 新版本可只用 KRaft(略)
bin/kafka-server-start.sh config/server.properties
bin/kafka-topics.sh \
--create \
--topic test \
--bootstrap-server localhost:9092 \
--partitions 1 \
--replication-factor 1
bin/kafka-console-consumer.sh \
--topic test \
--bootstrap-server localhost:9092 \
--from-beginning
✅ 此后只要在另一个终端生产消息,这里就能看到。
bin/kafka-console-producer.sh \
--topic test \
--bootstrap-server localhost:9092
输入:
hello kafka
常用参数(consumer.properties):
bootstrap.servers=localhost:9092
group.id=test-group
auto.offset.reset=earliest
enable.auto.commit=true
| 参数 | 说明 |
|---|---|
| group.id | 消费者组(非常重要) |
| auto.offset.reset | earliest / latest |
| enable.auto.commit | 是否自动提交 offset |
启动:
bin/kafka-console-consumer.sh \
--topic test \
--bootstrap-server localhost:9092 \
--consumer.config config/consumer.properties
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.6.1</version>
</dependency>
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "test-group");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("auto.offset.reset", "earliest");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(List.of("test"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> r : records) {
System.out.println(r.value());
}
}
pip install kafka-python
from kafka import KafkaConsumer
consumer = KafkaConsumer(
'test',
bootstrap_servers='localhost:9092',
group_id='test-group',
auto_offset_reset='earliest'
)
for msg in consumer:
print(msg.value.decode())
group.idauto.offset.reset=earliestsudo ufw allow 9092
如果你告诉我:
我可以给你 更精确的配置方案。