下面按 Debian 系统上使用 Kafka 消费者 的常见场景,给你一套清晰配置步骤(以官方 Kafka 二进制包为例,不依赖 Docker)。
Kafka 需要 Java:
sudo apt update
sudo apt install -y default-jdk
java -version
建议 JDK 11 或 17(Kafka 3.x 推荐)
cd /opt
sudo wget https://downloads.apache.org/kafka/3.7.0/kafka_2.13-3.7.0.tgz
sudo tar -xzf kafka_2.13-3.7.0.tgz
sudo mv kafka_2.13-3.7.0 kafka
cd /opt/kafka
bin/zookeeper-server-start.sh config/zookeeper.properties
bin/kafka-server-start.sh config/server.properties
Kafka 3.3+ 支持 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
✅ 此时你已经是“Kafka 消费者”了
常用配置文件(可新建 consumer.properties):
bootstrap.servers=localhost:9092
group.id=my-consumer-group
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
auto.offset.reset=earliest
enable.auto.commit=true
使用方式:
bin/kafka-console-consumer.sh \
--topic test-topic \
--consumer.config consumer.properties
| 参数 | 说明 |
|---|---|
| group.id | 消费者组(同组分担消费) |
| auto.offset.reset | earliest / latest |
| enable.auto.commit | 是否自动提交 offset |
| max.poll.records | 单次拉取数量 |
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.7.0</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");
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.println(r.value());
}
}
运行:
java -cp your.jar Main
nohup bin/kafka-server-start.sh config/server.properties > kafka.log 2>&1 &
或用 systemd(推荐)
listeners 和防火墙auto.offset.reset=earliestenable.auto.commit,手动提交如果你说的是:
告诉我,我直接给你对应配置。