在 Ubuntu 上配置 Kafka 消费者,通常分为 命令行方式(测试用) 和 程序方式(Java / Python 等)。下面从安装到配置一步一步说明。
# 安装 Java
sudo apt update
sudo apt install openjdk-11-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
# 启动 Kafka
bin/kafka-server-start.sh config/server.properties
(Kafka 3.x 也可用 KRaft 模式,不需要 Zookeeper)
bin/kafka-topics.sh \
--create \
--topic test-topic \
--bootstrap-server localhost:9092 \
--partitions 1 \
--replication-factor 1
bin/kafka-console-consumer.sh \
--topic test-topic \
--bootstrap-server localhost:9092 \
--from-beginning
参数说明:
--from-beginning:从最早消息开始消费Kafka 消费者主要配置项(命令行或代码中通用):
| 配置项 | 说明 |
|---|---|
bootstrap.servers |
Kafka 地址 |
group.id |
消费者组(非常重要) |
key.deserializer |
key 反序列化 |
value.deserializer |
value 反序列化 |
auto.offset.reset |
earliest / latest |
enable.auto.commit |
是否自动提交 offset |
bin/kafka-console-consumer.sh \
--topic test-topic \
--bootstrap-server localhost:9092 \
--group my-group
查看消费进度:
bin/kafka-consumer-groups.sh \
--describe \
--group my-group \
--bootstrap-server localhost:9092
<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("enable.auto.commit", "true");
props.put("auto.offset.reset", "earliest");
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(List.of("test-topic"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> r : records) {
System.out.printf("offset=%d, key=%s, value=%s%n",
r.offset(), r.key(), r.value());
}
}
pip install kafka-python
from kafka import KafkaConsumer
consumer = KafkaConsumer(
'test-topic',
bootstrap_servers='localhost:9092',
group_id='py-group',
auto_offset_reset='earliest'
)
for msg in consumer:
print(msg.value.decode())
bootstrap-server 是否正确group.idauto.offset.resetenable.auto.commitcommitSync()sudo ufw allow 9092
如果你告诉我:
我可以给你更精准的配置示例。