在 Debian 上连接 Kafka 集群,通常分三种场景来说:命令行连接、Java 客户端连接、其他语言客户端连接。下面以最常见的 Kafka 2.x / 3.x + Zookeeper 或 KRaft 为例。
连接 Kafka 集群你需要知道:
broker1:9092,broker2:9092,broker3:9092
telnet broker1 9092
sudo apt update
sudo apt install -y default-jre
下载 Kafka:
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
bin/kafka-topics.sh \
--bootstrap-server broker1:9092,broker2:9092,broker3:9092 \
--list
bin/kafka-console-producer.sh \
--bootstrap-server broker1:9092,broker2:9092,broker3:9092 \
--topic test
bin/kafka-console-consumer.sh \
--bootstrap-server broker1:9092,broker2:9092,broker3:9092 \
--topic test \
--from-beginning
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.6.1</version>
</dependency>
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("test", "hello kafka"));
producer.close();
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3: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"));
while (true) {
consumer.poll(Duration.ofMillis(100))
.forEach(r -> System.out.println(r.value()));
}
bootstrap.servers=broker1:9093,broker2:9093
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
ssl.truststore.location=/etc/kafka/client.truststore.jks
ssl.truststore.password=123456
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required \
username="user" \
password="pass";
# 检查端口
telnet broker1 9092
# 检查 hosts
cat /etc/hosts
LEADER_NOT_AVAILABLEbin/kafka-topics.sh --bootstrap-server ... --describe --topic test
sudo ufw allow 9092
核心就一句话:
Kafka 客户端只需要
bootstrap.servers,剩下的交给集群自己路由。
如果你愿意,可以告诉我:
我可以直接给你一份可复制的配置。