在 Debian 上配置 Kafka 消费者,一般分几种场景:
下面按常见情况给你说明。
sudo apt update
sudo apt install -y default-jdk
java -version
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 / Kafka(老版本用 Zookeeper,新版本可用 KRaft):
# Zookeeper 模式
bin/zookeeper-server-start.sh config/zookeeper.properties &
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
常用参数:
--group my-group:指定消费者组--max-messages 10:消费 10 条后退出--property print.key=true:打印 key示例(带 group):
bin/kafka-console-consumer.sh \
--topic test \
--bootstrap-server localhost:9092 \
--group test-group \
--from-beginning
Kafka 也支持 consumer.properties
vi config/consumer.properties
示例内容:
bootstrap.servers=localhost:9092
group.id=test-group
enable.auto.commit=true
auto.commit.interval.ms=1000
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
bin/kafka-console-consumer.sh \
--topic test \
--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("enable.auto.commit", "true");
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"));
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());
}
}
运行:
java -cp your.jar Main
示例服务文件:
sudo vim /etc/systemd/system/kafka-consumer.service
[Unit]
Description=Kafka Consumer
After=network.target
[Service]
ExecStart=/opt/kafka/bin/kafka-console-consumer.sh \
--topic test \
--bootstrap-server localhost:9092 \
--group test-group
Restart=always
User=kafka
[Install]
WantedBy=multi-user.target
sudo systemctl daemon-reload
sudo systemctl enable kafka-consumer
sudo systemctl start kafka-consumer
listeners 和 advertised.listeners--from-beginning如果你告诉我:
我可以给你 更精确的 Debian 配置示例。