下面按Debian 系统上配置 Kafka 消费者来给你讲,分 基础环境、消费者配置、命令行消费者、Java 消费者示例 四部分。
sudo apt update
sudo apt install -y default-jdk
java -version
wget https://downloads.apache.org/kafka/3.7.0/kafka_2.13-3.7.0.tgz
tar -xzf kafka_2.13-3.7.0.tgz
cd kafka_2.13-3.7.0
# 启动 Zookeeper(老版本)
bin/zookeeper-server-start.sh config/zookeeper.properties &
# 启动 Kafka
bin/kafka-server-start.sh config/server.properties &
Kafka 3.x 也可用 KRaft(无 Zookeeper),但初学建议先用上面方式。
消费者主要配置项(consumer.properties):
bootstrap.servers=localhost:9092
group.id=test-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
| 参数 | 说明 |
|---|---|
| bootstrap.servers | Kafka 地址 |
| group.id | 消费组 |
| auto.offset.reset | earliest / latest |
| enable.auto.commit | 是否自动提交 offset |
| session.timeout.ms | 消费者超时 |
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-consumer.sh \
--topic test \
--bootstrap-server localhost:9092 \
--group my-group
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.7.0</version>
</dependency>
import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.*;
public class KafkaConsumerDemo {
public static void main(String[] args) {
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(Collections.singletonList("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());
}
}
}
}
server.propertieslisteners=PLAINTEXT://0.0.0.0:9092
advertised.listeners=PLAINTEXT://<服务器IP>:9092
auto.offset.reset=earliest如果你告诉我:
我可以直接给你一份可复制的配置文件。