一、消息队列概述
1,什么是消息队列
· 传统业务系统的问题
1,用户访问效率问题
假设:我们自研的电商平台,双11有一个爆品被10w人同时点击购买;
- 订单请求会根据网站开发的程序循序,逐次执行;
- 每个业务程序执行都有时间消耗;
- 等到所有的程序执行成功,才会已返回给用户下单成功的消息;
- 导致问题:
| 问题 | 描述 |
|---|---|
| 【网站访问速度问题】 | 用户下单后,需要等待所有系统程序执行完成,才能得到下单结果反馈; 每个系统的执行速度根据当前服务器性能所决定,高并发时可能有较强延迟; |
| 【服务器高负载问题】 | 在特殊消费日来临时,高并发需要虚增服务器来满足; 而这些新增的服务器,在非消费日又用不到(还必须准备着),导致企业资源浪费; |

2,服务器资源利用问题
假设:
- 每天晚上只有三个小时是网站访问高峰期,需要12台服务器才能满足负载需求;
- 而每天剩余的21个小时,都是低谷期,只需要2台服务器就能满足负载要求;
- 那么也就是说,企业80%的资源,是限制的,浪费的;

· 消息队列的作用
消息队列【Message Queue】:简称【MQ】
| 功能 | 解释说明 |
|---|---|
| 消息并发 | 异步执行消息,使得执行效率翻倍 |
| 削峰填谷 | 缓存消息,使得服务器资源可以最大化利用 |
1,消息并发
【消息并发】:
- 就是将本来多系统的顺序执行逻辑,变为多系统并行执行;
- 原本需要等待上一步完成才能执行下一步,使用并发消息策略可以多系统异步并行执行;

2,削峰填谷
【削峰填谷】:
- 就是将大量并发数据,以缓存的方式管理起来;
- 然后有序地,释放执行,减缓后端系统服务器的压力;
- 保证后端服务器在高并发时刻,可以有序地,安全的处理请求;而不用增加服务器资源;

2,常见的MQ应用
| 常见MQ | 吞吐量 | 延迟 | 消息顺序 | 适合场景 | 不适合场景 | 丢数据 |
|---|---|---|---|---|---|---|
| Kafka | 100w+/秒 | 2-100毫秒 | 分区内严格有序 | 文本日志 JSON数据 事假记录 监控指标 交易记录 用户行为 |
大文件:图片/视频) 复杂路由需求 极低延迟<=1毫秒的 |
几乎不会 |
| RabbitMQ | 1-5w/秒 | 0.1-1毫秒 | 无保证 | 微服务通信 任务队列 复杂路由消息 RPC调用 工作流编排 |
大数据量日志 海量消息重放 长期存储 |
适中 |
| RocketMQ | 10-50w/秒 | 1-5毫秒 | 严格有序 | 金融交易数据 电商订单流程 账单处理 需要事务的消息 |
只适合java | 极低 |
| AxtiveMQ(已过时) | <=1w/秒 | 10-100毫秒 | 无保证 | 传统企业应用 | 极高 |
3,kafka数据与集群介绍
| 名词 | 解释说明 |
|---|---|
| broker | kafka集群的每个节点服务器,就是一个broker,broker列表:就是指kafka的集群所有节点 |
| producer | 【生产者】:就是往kafka集群【写入数据】的角色 |
| consumer | 【消费者】:就是找kafka集群【索要数据】的角色 |
| consumer group | 【消费者组】:就是找kafka集群【索要数据】的一组角色 |
| topic | 【主题】:就是【生产者】和【消费者】写入和索要数据的逻辑单元;【类似ES中的索引】 |
| parition | 【分区】:【topic主题】内部划分的区域;【类似ES中的分片】 |
| offset | 【数据消费记录】:【消费者】索要数据的位置,避免重复消费(重复读取); |
| replica | 【副本】:分区内的备份副本;【类似ES中的副本分片】 注意:与ES不同,这里的副本数是指,主分区+副本备份的总和; |
# 在kafka集群中,有两个维度的Leader和Follower角色
- 节点级别:leader主节点,Follower从节点;
- 数据级别:leader主数据,Follower副本数据;

二、单点部署kafka
1,环境准备
· 机器准备
| 主机名 | ip | 配置要求 |
|---|---|---|
| kafka01 | 10.0.0.111 | 2核4G |
| kafka02 | 10.0.0.112 | 2核4G |
| kafka03 | 10.0.0.113 | 2核4G |
· 主机名hosts解析
cat >> /etc/hosts <<EOF
10.0.0.111 kafka01
10.0.0.112 kafka02
10.0.0.113 kafka03
EOF
2,下载安装包
· 下载kafka安装包
1,官网
https://kafka.apache.org/

2,进入下载页面
# 点击右上角【DOWNLOAD KAFKA】进入下来页面

3,选择版本下载
| 版本 | 上线时间 | 说明 |
|---|---|---|
| 【3.3.2】 | 2023年1月 | kafka脱离zookeeper的最早稳定版本; |
# 选择3.3.2版本的2进制安装包

# 等待下载完成

· 下载jdk11安装包
1,官网
# 进入官网:
https://www.oracle.com/java/technologies/downloads/archive/

2,下载JDK11版本
# 选择X64社区版

# 等待下载完成

3,上传安装包
· 上传kafka安装包
[root@kafka01 ~ ]# rz -E
[root@kafka01 ~ ]# ls -l
-rw-r--r-- 1 root root 106619987 Dec 11 14:45 kafka_2.13-3.3.2.tgz
· 上传JDK安装包
[root@kafka01 ~ ]# rz -E
[root@kafka01 ~ ]# ls -l
-rw-r--r-- 1 root root 168251647 Dec 12 00:06 jdk-11.0.28_linux-x64_bin.tar.gz
-rw-r--r-- 1 root root 106619987 Dec 11 14:45 kafka_2.13-3.3.2.tgz
4,创建安装目录
[root@kafka01 ~ ]# mkdir /tools
5,安装JDK
· 解压JDK到安装目录
[root@kafka01 ~ ]# tar xf jdk-11.0.28_linux-x64_bin.tar.gz -C /tools/
[root@kafka01 ~ ]# ls -l /tools/
drwxr-xr-x 9 root root 4096 Dec 12 00:12 jdk-11.0.28
· 配置环境变量
[root@kafka01 ~ ]# vim /etc/profile
......
export JAVA_HOME="/tools/jdk-11.0.28"
export PATH="$PATH:$JAVA_HOME/bin:$JAVA_HOME/jre/bin"
export CLASSPATH="$CLASSPATH:$JAVA_HOME/lib/:$JAVA_HOME/jre/lib:$JAVA_HOME/lib/tools.jar"
# 变量生效
[root@kafka01 ~ ]# source /etc/profile
# 查看java
[root@kafka01 ~ ]# java -version
java version "11.0.28" 2025-07-15 LTS
Java(TM) SE Runtime Environment 18.9 (build 11.0.28+12-LTS-279)
Java HotSpot(TM) 64-Bit Server VM 18.9 (build 11.0.28+12-LTS-279, mixed mode)
6,单点部署kafka
· 解压安装包
[root@kafka01 ~ ]# tar xf kafka_2.13-3.3.2.tgz -C /tools/
[root@kafka01 ~ ]# ls -l /tools/kafka_2.13-3.3.2/
total 64
drwxr-xr-x 3 root root 4096 Dec 21 2022 bin
drwxr-xr-x 3 root root 4096 Dec 21 2022 config
drwxr-xr-x 2 root root 4096 Dec 11 23:51 libs
-rw-r--r-- 1 root root 14844 Dec 21 2022 LICENSE
drwxr-xr-x 2 root root 4096 Dec 21 2022 licenses
-rw-r--r-- 1 root root 28184 Dec 21 2022 NOTICE
drwxr-xr-x 2 root root 4096 Dec 21 2022 site-docs
· 创建数据和证书存储目录
[root@kafka01 ~ ]# mkdir /tools/kafka_2.13-3.3.2/{data,certs}
[root@kafka01 ~ ]# ls -l /tools/kafka_2.13-3.3.2/
total 72
drwxr-xr-x 3 root root 4096 Dec 21 2022 bin
drwxr-xr-x 3 root root 4096 Dec 21 2022 config
drwxr-xr-x 2 root root 4096 Dec 12 04:47 certs
drwxr-xr-x 2 root root 4096 Dec 12 04:47 data
drwxr-xr-x 2 root root 4096 Dec 11 23:51 libs
-rw-r--r-- 1 root root 14844 Dec 21 2022 LICENSE
drwxr-xr-x 2 root root 4096 Dec 21 2022 licenses
-rw-r--r-- 1 root root 28184 Dec 21 2022 NOTICE
drwxr-xr-x 2 root root 4096 Dec 21 2022 site-docs
· 生成证书
1,生成CA根证书
# 进入证书目录
[root@kafka01 ~ ]# cd /tools/kafka_2.13-3.3.2/certs
# 生成ca的私钥和公钥
[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# openssl req -x509 -newkey rsa:4096 -sha256 -days 3650 -nodes \
-subj "/C=CN/ST=Beijing/L=Beijing/O=bakwite/CN=kafka Root CA" \
-keyout kafka-ca.key -out kafka-ca.crt
2,生成服务端证书
# 生成服务端私钥与请求文件
[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# openssl req -newkey rsa:4096 -nodes -sha256 \
-subj "/C=CN/ST=Beijing/L=Beijing/O=bakwite/CN=kafka-cluster" \
-addext "subjectAltName=DNS:kafka01,DNS:kafka02,DNS:kafka03,IP:10.0.0.111,IP:10.0.0.112,IP:10.0.0.113,IP:127.0.0.1" \
-keyout kafka-server.key -out kafka-server.csr
# 生成服务端公钥(集群通讯证书)
[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# openssl x509 -req -in kafka-server.csr -CA kafka-ca.crt -CAkey kafka-ca.key -CAcreateserial \
-days 365 -out kafka-server.crt \
-extfile <(printf "subjectAltName=DNS:kafka01,DNS:kafka02,DNS:kafka03,IP:10.0.0.111,IP:10.0.0.112,IP:10.0.0.113,IP:127.0.0.1")
3,服务端证书合并
[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# cat kafka-server.crt kafka-server.key > kafka-server.pem
查看所有证书
[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# ls -l
-rw-r--r-- 1 kafka kafka 2009 Dec 13 10:24 kafka-ca.crt
-rw------- 1 kafka kafka 3272 Dec 13 10:24 kafka-ca.key
-rw-r--r-- 1 kafka kafka 2069 Dec 13 10:25 kafka-server.crt
-rw-r--r-- 1 kafka kafka 1793 Dec 13 10:25 kafka-server.csr
-rw------- 1 kafka kafka 3272 Dec 13 10:25 kafka-server.key
-rw-r--r-- 1 kafka kafka 5341 Dec 13 10:25 kafka-server.pem
· 编辑配置文件
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
# ===== 基础配置 =====
# 节点角色
# 控制器节点 (Controller):仅参与集群管理,不处理客户端的生产消费请求。
# 代理节点 (Broker):仅负责存储数据和响应客户端请求,不参与集群管理。
process.roles=broker,controller
# 节点id(集群唯一)
node.id=1
# 总裁选举列表
# # 集群部署
# controller.quorum.voters=1@kafka01:9093,2@kafka02:9093,3@kafka03:9093
# # 单点部署
controller.quorum.voters=1@kafka01:9093
# ===== 网络配置 =====
# 【SASL_SSL://IP:9092】:客户端连接入口,(SASL_SSL表示证书加用户名密码)
listener.security.protocol.map=CONTROLLER:PLAINTEXT,SASL_SSL:SASL_SSL
listeners=SASL_SSL://10.0.0.111:9092,CONTROLLER://10.0.0.111:9093
# 对外公布的地址:客户端实际连接的地址
advertised.listeners=SASL_SSL://10.0.0.111:9092
# broker间通信使用的监听器名称
inter.broker.listener.name=SASL_SSL
# ===== 存储配置 =====
# 数据持久化存储目录,可以写多个,以【,】分割
log.dirs=/tools/kafka_2.13-3.3.2/data
# 默认topic数据保存的时间(毫秒)
log.retention.ms=604800000
# 默认topic数据保存的大小限制(byte)
log.segment.bytes=1073741824
# 消费者offset记录保存时间(单位:分钟[7天])
offsets.retention.minutes=10080
# ===== 集群配置 =====
# 禁止自动创建topic
auto.create.topics.enable=false
# 创建topic时的默认分区数,分区数决定并行度,通常按业务需求设置
# 单点部署设置为1,集群部署设置为【集群数】
num.partitions=1
# 默认副本因子:每个分区的副本数,通常设为3,保证高可用
# 单点部署设置为1,集群部署设置为【集群数】
default.replication.factor=1
# 最小同步副本数,当生产者来确认时,需要这么多副本确认
# 单点部署设置为1,集群部署设置为【集群数-1】
min.insync.replicas=1
# 【__consumer_offsets】主题的副本数,存储消费者偏移量,很重要;
# 单点部署设置为1,集群部署设置为【集群数】
offsets.topic.replication.factor=1
# 事务状态日志的副本数;(如果使用Kafka事务才起作用)
# 单点部署设置为1,集群部署设置为【集群数】
transaction.state.log.replication.factor=1
# 事务状态日志的最小同步副本数
# 单点部署设置为1,集群部署设置为【集群数-1】
transaction.state.log.min.isr=1
# ===== 控制器配置 =====
# Controller监听器名称,必须与listeners中的CONTROLLER对应
controller.listener.names=CONTROLLER
# ===== 性能优化 =====
# 线程数设置(CPU核心数的2倍)
num.network.threads=4
# IO线程数:处理磁盘读写(磁盘数*8)
num.io.threads=8
# socket发送缓冲区大小:1MB
socket.send.buffer.bytes=1024000
# socket接收缓冲区大小:1MB
socket.receive.buffer.bytes=1024000
# 最大请求大小:100MB
socket.request.max.bytes=104857600
# ===== 副本配置 =====
# 主分区宕机时,不允许未同步完成同步的副本分区成为leader分区
unclean.leader.election.enable=false
# 是否自动平衡leader分布
auto.leader.rebalance.enable=true
# 副本滞后时间阈值:30秒(超过则认为副本失效)
replica.lag.time.max.ms=30000
# ===== 安全认证配置 =====
# SSL配置
# 指定PEM格式
ssl.keystore.type=PEM
ssl.keystore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
# 关闭客户端双向SSL认证(开启双向认证:required)
ssl.client.auth=none
# 明确设为空,表示不验证主机名
ssl.endpoint.identification.algorithm=
# 在配置文件中添加信任库配置[自己链接自己时需要信任证书]
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-ca.crt
# SASL配置(用户密码认证)启用的SASL机制:
# # 【PLAIN】(用户名密码明文传输)【SCRAM-SHA-256】(更安全)
sasl.enabled.mechanisms=PLAIN
# broker间通信使用的SASL机制
sasl.mechanism.inter.broker.protocol=PLAIN
· 创建kafka用户名密码文件
写法:[user_用户名="密码"]
说明:
- 【KafkaServer用户】:是broker节点接收链接请求时的认证用户;
- 【username/password】:表示节点之间通讯验证的用户;
- 【user_用户名="密码"】:表示客户端的链接验证用户;
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/kafka_users.conf
KafkaServer {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="broker"
password="123456"
user_admin="123456"
user_filebeat="123456"
user_logstash="123456"
user_broker="123456";
};
· 编辑启动脚本
将用户文件加入到启动文件中,使其生效;
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/bin/kafka-server-start.sh
export KAFKA_OPTS="-Djava.security.auth.login.config=/tools/kafka_2.13-3.3.2/config/kafka_users.conf"
......
· 生成集群ID格式化存储目录
注意:单点部署也需要格式化存储目录
1,生成集群ID
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-storage.sh random-uuid
zP0Jy27XTyKo090WVkFqUQ
2,格式化存储目录
必须格式化,否则存储会有问题
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-storage.sh format \
-t zP0Jy27XTyKo090WVkFqUQ \
-c /tools/kafka_2.13-3.3.2/config/server.properties
Formatting /tools/kafka_2.13-3.3.2/data with metadata.version 3.3-IV3.
· 前台启动kafka测试(可以不前台启动)
1,前台启动
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-server-start.sh /tools/kafka_2.13-3.3.2/config/server.properties
2,查看启动端口
[root@kafka01 ~ ]# ss -tnulp
tcp LISTEN 0 50 [::ffff:10.0.0.111]:9092 *:* users:(("java",pid=124457,fd=144))
tcp LISTEN 0 50 [::ffff:10.0.0.111]:9093 *:* users:(("java",pid=124457,fd=124))
· 配置system启动
先【ctrl + c】关闭kafka前台运行
1,创建kafka用户
[root@kafka01 ~ ]# useradd -M -s /sbin/nologin kafka
2,授权kafka安装目录
[root@kafka01 ~ ]# chown -R kafka.kafka /tools/kafka_2.13-3.3.2/
3,编辑system文件
[root@kafka01 ~ ]# vim /lib/systemd/system/kafka.service
[Unit]
Description=Apache Kafka Server
After=network.target
[Service]
Type=simple
User=kafka
Group=kafka
Environment="JAVA_HOME=/tools/jdk-11.0.28"
Environment="KAFKA_HEAP_OPTS=-Xmx1g -Xms1g"
ExecStart=/tools/kafka_2.13-3.3.2/bin/kafka-server-start.sh /tools/kafka_2.13-3.3.2/config/server.properties
ExecStop=/tools/kafka_2.13-3.3.2/bin/kafka-server-stop.sh
Restart=on-failure
RestartSec=10
LimitNOFILE=65536
[Install]
WantedBy=multi-user.target
4,system启动kafka
[root@kafka01 ~ ]# systemctl daemon-reload
[root@kafka01 ~ ]# systemctl enable --now kafka.service
三、集群部署kafka
1,分发安装目录文件
· 关闭单点kafka
[root@kafka01 ~ ]# systemctl stop kafka.service
· 分发安装目录给其他节点
[root@kafka01 ~ ]# scp -r /tools 10.0.0.112:/
[root@kafka01 ~ ]# scp -r /tools 10.0.0.113:/
· 分发全局变量文件
[root@kafka01 ~ ]# scp /etc/profile 10.0.0.112:/etc/
[root@kafka01 ~ ]# scp /etc/profile 10.0.0.113:/etc/
# 拷贝后执行生效
[root@kafka02 ~ ]# source /etc/profile
[root@kafka03 ~ ]# source /etc/profile
· 分发system启动文件
[root@kafka01 ~ ]# scp /lib/systemd/system/kafka.service 10.0.0.112:/lib/systemd/system/
[root@kafka01 ~ ]# scp /lib/systemd/system/kafka.service 10.0.0.113:/lib/systemd/system/
2,其他节点创建kafka用户
[root@kafka02 ~ ]# useradd -M -s /sbin/nologin kafka
[root@kafka03 ~ ]# useradd -M -s /sbin/nologin kafka
3,其他节点授权kafka目录
[root@kafka02 ~ ]# chown -R kafka.kafka /tools/kafka_2.13-3.3.2/
[root@kafka03 ~ ]# chown -R kafka.kafka /tools/kafka_2.13-3.3.2/
4,编辑三个节点配置文件
· kafka01节点
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
# ===== 基础配置 =====
# 节点角色
process.roles=broker,controller
# 节点id(集群唯一)
node.id=1
# 总裁选举列表
# # 集群部署
controller.quorum.voters=1@kafka01:9093,2@kafka02:9093,3@kafka03:9093
# # 单点部署
#controller.quorum.voters=1@kafka01:9093
# ===== 网络配置 =====
# 【SASL_SSL://IP:9092】:客户端连接入口,(SASL_SSL表示证书加用户名密码)
listener.security.protocol.map=CONTROLLER:PLAINTEXT,SASL_SSL:SASL_SSL
listeners=SASL_SSL://10.0.0.111:9092,CONTROLLER://10.0.0.111:9093
# 对外公布的地址:客户端实际连接的地址
advertised.listeners=SASL_SSL://10.0.0.111:9092
# broker间通信使用的监听器名称
inter.broker.listener.name=SASL_SSL
# ===== 存储配置 =====
# 数据持久化存储目录,可以写多个,以【,】分割
log.dirs=/tools/kafka_2.13-3.3.2/data
# 默认topic数据保存的时间(毫秒)
log.retention.ms=604800000
# 默认topic数据保存的最大大小(byte);
log.segment.bytes=1073741824
# ===== 集群配置 =====
# 关闭自动创建topic
auto.create.topics.enable=false
# 创建topic时的默认分区数,分区数决定并行度,通常按业务需求设置
# 单点部署设置为1,集群部署设置为【集群数】
num.partitions=3
# 默认副本因子:每个topic分区的副本数,通常设为3,保证高可用
# 单点部署设置为1,集群部署设置为【集群数】
default.replication.factor=3
# 最小同步副本数,当生产者写入数据时,需要等待起码多少个副本写入完成,才算做写入数据成功?
# 单点部署设置为1,集群部署设置为【集群数-1】
min.insync.replicas=2
# 【__consumer_offsets】主题的副本数,存储消费者偏移量,很重要;
# 单点部署设置为1,集群部署设置为【集群数】
offsets.topic.replication.factor=3
# 事务状态日志的副本数;(如果使用Kafka事务才起作用)
# 单点部署设置为1,集群部署设置为【集群数】
transaction.state.log.replication.factor=3
# 事务状态日志的最小同步副本数
# 单点部署设置为1,集群部署设置为【集群数-1】
transaction.state.log.min.isr=2
# 消费者offset记录保存时间(分钟)
offsets.retention.minutes=10080
# ===== 控制器配置 =====
# Controller监听器名称,必须与listeners中的CONTROLLER对应
controller.listener.names=CONTROLLER
# ===== 性能优化 =====
# 线程数设置(CPU核心数的2倍)
num.network.threads=4
# IO线程数:处理磁盘读写(磁盘数*8)
num.io.threads=8
# socket发送缓冲区大小:1MB
socket.send.buffer.bytes=1024000
# socket接收缓冲区大小:1MB
socket.receive.buffer.bytes=1024000
# 最大请求大小:100MB
socket.request.max.bytes=104857600
# ===== 副本配置 =====
# 主分区宕机时,不允许未同步完成同步的副本分区成为leader分区
unclean.leader.election.enable=false
# 是否自动平衡leader分布
auto.leader.rebalance.enable=true
# 副本滞后时间阈值:30秒(超过则认为副本失效)
replica.lag.time.max.ms=30000
# ===== 安全认证配置 =====
# SSL配置
# 指定PEM格式
ssl.keystore.type=PEM
ssl.keystore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
# 关闭客户端双向SSL认证
ssl.client.auth=none
# 明确设为空,表示不验证主机名
ssl.endpoint.identification.algorithm=
# 在配置文件中添加信任库配置[自己链接自己时需要信任证书]
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
# SASL配置(用户密码认证)
# 启用的SASL机制:
# # 【PLAIN】(用户名密码明文传输)【SCRAM-SHA-256】(更安全)
sasl.enabled.mechanisms=PLAIN
# broker间通信使用的SASL机制
sasl.mechanism.inter.broker.protocol=PLAIN
· kafka02节点
[root@kafka02 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
process.roles=broker,controller
# 节点id(集群唯一):这里每个节点都要不一样
node.id=2
# 总裁选举列表
# # 集群部署
controller.quorum.voters=1@kafka01:9093,2@kafka02:9093,3@kafka03:9093
# # 单点部署
#controller.quorum.voters=1@kafka01:9093
# 这里ip地址改为节点的IP
# 【SASL_SSL://IP:9092】:客户端连接入口,(SASL_SSL表示证书加用户名密码)
listener.security.protocol.map=CONTROLLER:PLAINTEXT,SASL_SSL:SASL_SSL
listeners=SASL_SSL://10.0.0.112:9092,CONTROLLER://10.0.0.112:9093
# 对外公布的地址:客户端实际连接的地址
advertised.listeners=SASL_SSL://10.0.0.112:9092
# broker间通信使用的监听器名称
inter.broker.listener.name=SASL_SSL
# 以下:【所有节点都一样】
log.dirs=/tools/kafka_2.13-3.3.2/data
# 默认topic数据保存的时间(毫秒)
log.retention.ms=604800000
log.segment.bytes=1073741824
auto.create.topics.enable=false
num.partitions=3
default.replication.factor=3
min.insync.replicas=2
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
offsets.retention.minutes=10080
controller.listener.names=CONTROLLER
num.network.threads=4
num.io.threads=8
socket.send.buffer.bytes=1024000
socket.receive.buffer.bytes=1024000
socket.request.max.bytes=104857600
unclean.leader.election.enable=false
auto.leader.rebalance.enable=true
replica.lag.time.max.ms=30000
ssl.keystore.type=PEM
ssl.keystore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
ssl.client.auth=none
ssl.endpoint.identification.algorithm=
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN
· kafka03节点
[root@kafka03 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
process.roles=broker,controller
# 节点id(集群唯一):这里每个节点都要不一样
node.id=3
# 总裁选举列表
# # 集群部署
controller.quorum.voters=1@kafka01:9093,2@kafka02:9093,3@kafka03:9093
# # 单点部署
#controller.quorum.voters=1@kafka01:9093
# 这里ip地址改为节点的IP
# 【SASL_SSL://IP:9092】:客户端连接入口,(SASL_SSL表示证书加用户名密码)
listener.security.protocol.map=CONTROLLER:PLAINTEXT,SASL_SSL:SASL_SSL
listeners=SASL_SSL://10.0.0.113:9092,CONTROLLER://10.0.0.113:9093
# 对外公布的地址:客户端实际连接的地址
advertised.listeners=SASL_SSL://10.0.0.113:9092
# broker间通信使用的监听器名称
inter.broker.listener.name=SASL_SSL
# 以下:【所有节点都一样】
log.dirs=/tools/kafka_2.13-3.3.2/data
# 默认topic数据保存的时间(毫秒)
log.retention.ms=604800000
log.segment.bytes=1073741824
auto.create.topics.enable=false
num.partitions=3
default.replication.factor=3
min.insync.replicas=2
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
offsets.retention.minutes=10080
controller.listener.names=CONTROLLER
num.network.threads=4
num.io.threads=8
socket.send.buffer.bytes=1024000
socket.receive.buffer.bytes=1024000
socket.request.max.bytes=104857600
unclean.leader.election.enable=false
auto.leader.rebalance.enable=true
replica.lag.time.max.ms=30000
ssl.keystore.type=PEM
ssl.keystore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
ssl.client.auth=none
ssl.endpoint.identification.algorithm=
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN
5,清理历史数据
之前单点部署时,产生的集群元数据;
[root@kafka01 ~ ]# rm -rf /tools/kafka_2.13-3.3.2/data/*
[root@kafka02 ~ ]# rm -rf /tools/kafka_2.13-3.3.2/data/*
[root@kafka03 ~ ]# rm -rf /tools/kafka_2.13-3.3.2/data/*
6,格式化目录
· kafka01生成集群ID
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-storage.sh random-uuid
u7LyLt7yRj-EaS8g2yh32g
· 所有节点初始化数据目录
/tools/kafka_2.13-3.3.2/bin/kafka-storage.sh format \
-t u7LyLt7yRj-EaS8g2yh32g \
-c /tools/kafka_2.13-3.3.2/config/server.properties
7,启动集群
· 所有节点system启动
systemctl daemon-reload
systemctl enable kafka.service
systemctl start kafka.service
· 节点查看
[root@kafka01 ~ ]# ss -tnulp
......
tcp LISTEN 0 50 [::ffff:10.0.0.111]:9092 *:*
tcp LISTEN 0 50 [::ffff:10.0.0.111]:9093 *:*
8,创建客户端用户(文件)
[root@kafka01 ~ ]# vim login-kafka.conf
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required \
username="filebeat" \
password="123456";
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
9,查看验证集群
· 查看集群元数据
| 显示字段 | 解释说明 |
|---|---|
| ClusterId: | 集群的ID(部署时,格式化目录生成的那个个) |
| LeaderId: | 主节点是哪个ID |
| LeaderEpoch: | 集群的leader更换了多少次了? |
| HighWatermark: | 【高水位线】:已经提交的元数据的日志偏移量 |
| MaxFollowerLag: | 【最大落后偏移量】:同步数据落后了多少偏移量 |
| MaxFollowerLagTimeMs: | 【最大落后时间】:同步数据落后了多长时间 |
| CurrentVoters: | 当前参与选举的broker(集群节点)列表 |
| CurrentObservers: | 当前不参与选举的(观察者)broker列表 |
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-metadata-quorum.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
describe --status
ClusterId: cuR3KFFBR8S9OlmZAduISw
LeaderId: 1
LeaderEpoch: 371
HighWatermark: 1506
MaxFollowerLag: 0
MaxFollowerLagTimeMs: 0
CurrentVoters: [1,2,3]
CurrentObservers: []
· 查看集群节点详细信息
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-broker-api-versions.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf
10.0.0.111:9092 (id: 1 rack: null) -> (
Produce(0): 0 to 9 [usable: 9],
Fetch(1): 0 to 13 [usable: 13],
ListOffsets(2): 0 to 7 [usable: 7],
Metadata(3): 0 to 12 [usable: 12],
LeaderAndIsr(4): UNSUPPORTED,
StopReplica(5): UNSUPPORTED,
UpdateMetadata(6): UNSUPPORTED,
ControlledShutdown(7): UNSUPPORTED,
OffsetCommit(8): 0 to 8 [usable: 8],
OffsetFetch(9): 0 to 8 [usable: 8],
FindCoordinator(10): 0 to 4 [usable: 4],
JoinGroup(11): 0 to 9 [usable: 9],
Heartbeat(12): 0 to 4 [usable: 4],
LeaveGroup(13): 0 to 5 [usable: 5],
SyncGroup(14): 0 to 5 [usable: 5],
DescribeGroups(15): 0 to 5 [usable: 5],
ListGroups(16): 0 to 4 [usable: 4],
SaslHandshake(17): 0 to 1 [usable: 1],
ApiVersions(18): 0 to 3 [usable: 3],
CreateTopics(19): 0 to 7 [usable: 7],
DeleteTopics(20): 0 to 6 [usable: 6],
DeleteRecords(21): 0 to 2 [usable: 2],
InitProducerId(22): 0 to 4 [usable: 4],
OffsetForLeaderEpoch(23): 0 to 4 [usable: 4],
AddPartitionsToTxn(24): 0 to 3 [usable: 3],
AddOffsetsToTxn(25): 0 to 3 [usable: 3],
EndTxn(26): 0 to 3 [usable: 3],
WriteTxnMarkers(27): 0 to 1 [usable: 1],
TxnOffsetCommit(28): 0 to 3 [usable: 3],
DescribeAcls(29): 0 to 3 [usable: 3],
CreateAcls(30): 0 to 3 [usable: 3],
DeleteAcls(31): 0 to 3 [usable: 3],
DescribeConfigs(32): 0 to 4 [usable: 4],
AlterConfigs(33): 0 to 2 [usable: 2],
AlterReplicaLogDirs(34): 0 to 2 [usable: 2],
DescribeLogDirs(35): 0 to 4 [usable: 4],
SaslAuthenticate(36): 0 to 2 [usable: 2],
CreatePartitions(37): 0 to 3 [usable: 3],
CreateDelegationToken(38): UNSUPPORTED,
RenewDelegationToken(39): UNSUPPORTED,
ExpireDelegationToken(40): UNSUPPORTED,
DescribeDelegationToken(41): UNSUPPORTED,
DeleteGroups(42): 0 to 2 [usable: 2],
ElectLeaders(43): 0 to 2 [usable: 2],
IncrementalAlterConfigs(44): 0 to 1 [usable: 1],
AlterPartitionReassignments(45): 0 [usable: 0],
ListPartitionReassignments(46): 0 [usable: 0],
OffsetDelete(47): 0 [usable: 0],
DescribeClientQuotas(48): 0 to 1 [usable: 1],
AlterClientQuotas(49): 0 to 1 [usable: 1],
DescribeUserScramCredentials(50): UNSUPPORTED,
AlterUserScramCredentials(51): UNSUPPORTED,
DescribeQuorum(55): 0 to 1 [usable: 1],
AlterPartition(56): UNSUPPORTED,
UpdateFeatures(57): 0 to 1 [usable: 1],
DescribeCluster(60): 0 [usable: 0],
DescribeProducers(61): 0 [usable: 0],
UnregisterBroker(64): 0 [usable: 0],
DescribeTransactions(65): 0 [usable: 0],
ListTransactions(66): 0 [usable: 0],
AllocateProducerIds(67): UNSUPPORTED
)
10.0.0.112:9092 (id: 2 rack: null) -> (
Produce(0): 0 to 9 [usable: 9],
Fetch(1): 0 to 13 [usable: 13],
ListOffsets(2): 0 to 7 [usable: 7],
Metadata(3): 0 to 12 [usable: 12],
LeaderAndIsr(4): UNSUPPORTED,
StopReplica(5): UNSUPPORTED,
UpdateMetadata(6): UNSUPPORTED,
ControlledShutdown(7): UNSUPPORTED,
OffsetCommit(8): 0 to 8 [usable: 8],
OffsetFetch(9): 0 to 8 [usable: 8],
FindCoordinator(10): 0 to 4 [usable: 4],
JoinGroup(11): 0 to 9 [usable: 9],
Heartbeat(12): 0 to 4 [usable: 4],
LeaveGroup(13): 0 to 5 [usable: 5],
SyncGroup(14): 0 to 5 [usable: 5],
DescribeGroups(15): 0 to 5 [usable: 5],
ListGroups(16): 0 to 4 [usable: 4],
SaslHandshake(17): 0 to 1 [usable: 1],
ApiVersions(18): 0 to 3 [usable: 3],
CreateTopics(19): 0 to 7 [usable: 7],
DeleteTopics(20): 0 to 6 [usable: 6],
DeleteRecords(21): 0 to 2 [usable: 2],
InitProducerId(22): 0 to 4 [usable: 4],
OffsetForLeaderEpoch(23): 0 to 4 [usable: 4],
AddPartitionsToTxn(24): 0 to 3 [usable: 3],
AddOffsetsToTxn(25): 0 to 3 [usable: 3],
EndTxn(26): 0 to 3 [usable: 3],
WriteTxnMarkers(27): 0 to 1 [usable: 1],
TxnOffsetCommit(28): 0 to 3 [usable: 3],
DescribeAcls(29): 0 to 3 [usable: 3],
CreateAcls(30): 0 to 3 [usable: 3],
DeleteAcls(31): 0 to 3 [usable: 3],
DescribeConfigs(32): 0 to 4 [usable: 4],
AlterConfigs(33): 0 to 2 [usable: 2],
AlterReplicaLogDirs(34): 0 to 2 [usable: 2],
DescribeLogDirs(35): 0 to 4 [usable: 4],
SaslAuthenticate(36): 0 to 2 [usable: 2],
CreatePartitions(37): 0 to 3 [usable: 3],
CreateDelegationToken(38): UNSUPPORTED,
RenewDelegationToken(39): UNSUPPORTED,
ExpireDelegationToken(40): UNSUPPORTED,
DescribeDelegationToken(41): UNSUPPORTED,
DeleteGroups(42): 0 to 2 [usable: 2],
ElectLeaders(43): 0 to 2 [usable: 2],
IncrementalAlterConfigs(44): 0 to 1 [usable: 1],
AlterPartitionReassignments(45): 0 [usable: 0],
ListPartitionReassignments(46): 0 [usable: 0],
OffsetDelete(47): 0 [usable: 0],
DescribeClientQuotas(48): 0 to 1 [usable: 1],
AlterClientQuotas(49): 0 to 1 [usable: 1],
DescribeUserScramCredentials(50): UNSUPPORTED,
AlterUserScramCredentials(51): UNSUPPORTED,
DescribeQuorum(55): 0 to 1 [usable: 1],
AlterPartition(56): UNSUPPORTED,
UpdateFeatures(57): 0 to 1 [usable: 1],
DescribeCluster(60): 0 [usable: 0],
DescribeProducers(61): 0 [usable: 0],
UnregisterBroker(64): 0 [usable: 0],
DescribeTransactions(65): 0 [usable: 0],
ListTransactions(66): 0 [usable: 0],
AllocateProducerIds(67): UNSUPPORTED
)
10.0.0.113:9092 (id: 3 rack: null) -> (
Produce(0): 0 to 9 [usable: 9],
Fetch(1): 0 to 13 [usable: 13],
ListOffsets(2): 0 to 7 [usable: 7],
Metadata(3): 0 to 12 [usable: 12],
LeaderAndIsr(4): UNSUPPORTED,
StopReplica(5): UNSUPPORTED,
UpdateMetadata(6): UNSUPPORTED,
ControlledShutdown(7): UNSUPPORTED,
OffsetCommit(8): 0 to 8 [usable: 8],
OffsetFetch(9): 0 to 8 [usable: 8],
FindCoordinator(10): 0 to 4 [usable: 4],
JoinGroup(11): 0 to 9 [usable: 9],
Heartbeat(12): 0 to 4 [usable: 4],
LeaveGroup(13): 0 to 5 [usable: 5],
SyncGroup(14): 0 to 5 [usable: 5],
DescribeGroups(15): 0 to 5 [usable: 5],
ListGroups(16): 0 to 4 [usable: 4],
SaslHandshake(17): 0 to 1 [usable: 1],
ApiVersions(18): 0 to 3 [usable: 3],
CreateTopics(19): 0 to 7 [usable: 7],
DeleteTopics(20): 0 to 6 [usable: 6],
DeleteRecords(21): 0 to 2 [usable: 2],
InitProducerId(22): 0 to 4 [usable: 4],
OffsetForLeaderEpoch(23): 0 to 4 [usable: 4],
AddPartitionsToTxn(24): 0 to 3 [usable: 3],
AddOffsetsToTxn(25): 0 to 3 [usable: 3],
EndTxn(26): 0 to 3 [usable: 3],
WriteTxnMarkers(27): 0 to 1 [usable: 1],
TxnOffsetCommit(28): 0 to 3 [usable: 3],
DescribeAcls(29): 0 to 3 [usable: 3],
CreateAcls(30): 0 to 3 [usable: 3],
DeleteAcls(31): 0 to 3 [usable: 3],
DescribeConfigs(32): 0 to 4 [usable: 4],
AlterConfigs(33): 0 to 2 [usable: 2],
AlterReplicaLogDirs(34): 0 to 2 [usable: 2],
DescribeLogDirs(35): 0 to 4 [usable: 4],
SaslAuthenticate(36): 0 to 2 [usable: 2],
CreatePartitions(37): 0 to 3 [usable: 3],
CreateDelegationToken(38): UNSUPPORTED,
RenewDelegationToken(39): UNSUPPORTED,
ExpireDelegationToken(40): UNSUPPORTED,
DescribeDelegationToken(41): UNSUPPORTED,
DeleteGroups(42): 0 to 2 [usable: 2],
ElectLeaders(43): 0 to 2 [usable: 2],
IncrementalAlterConfigs(44): 0 to 1 [usable: 1],
AlterPartitionReassignments(45): 0 [usable: 0],
ListPartitionReassignments(46): 0 [usable: 0],
OffsetDelete(47): 0 [usable: 0],
DescribeClientQuotas(48): 0 to 1 [usable: 1],
AlterClientQuotas(49): 0 to 1 [usable: 1],
DescribeUserScramCredentials(50): UNSUPPORTED,
AlterUserScramCredentials(51): UNSUPPORTED,
DescribeQuorum(55): 0 to 1 [usable: 1],
AlterPartition(56): UNSUPPORTED,
UpdateFeatures(57): 0 to 1 [usable: 1],
DescribeCluster(60): 0 [usable: 0],
DescribeProducers(61): 0 [usable: 0],
UnregisterBroker(64): 0 [usable: 0],
DescribeTransactions(65): 0 [usable: 0],
ListTransactions(66): 0 [usable: 0],
AllocateProducerIds(67): UNSUPPORTED
)
· 查看集群健康状态
| 显示字段 | 含义说明 |
|---|---|
| NodeId | 集群的ID |
| LogEndOffset | 元数据日志的结束偏移量(现在有2039条日志数据) |
| Lag | 副本落后于Leader的偏移量(表示数据完全同步) |
| LastFetchTimestamp | 最后一次拉取数据的时间戳 |
| LastCaughtUpTimestamp | 最后一次追上Leader的时间戳 |
| Status | 节点的角色 |
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-metadata-quorum.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
describe --replication
NodeId LogEndOffset Lag LastFetchTimestamp LastCaughtUpTimestamp Status
1 2039 0 1765622013103 1765622013103 Leader
2 2039 0 1765622012819 1765622012819 Follower
3 2039 0 1765622012819 1765622012819 Follower
四、kafka数据流管理
1,topic管理
先查看topic列表:当前什么都没有
[root@kafka01 ~ ]# cd /tools/kafka_2.13-3.3.2/
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--list
· 创建topic
| 参数 | 说明 |
|---|---|
| --create | 表示对topic的动作为:创建 |
| --bootstrap-server | 指定创建的节点 |
| --command-config | 指定使用的用户身份文件 |
| --partitions 3 | 指定分区数(不指定会使用配置文件中的默认参数值) |
| --replication-factor 3 | 指定副本数(不指定会使用配置文件中的默认参数值) 注意:kafka的分区副本与ES的概念不同,是本分数加上原数据的总数为副本数 |
[root@kafka01 ~ ]# cd /tools/kafka_2.13-3.3.2/ #login-kafka.conf必须在当前目录
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--create \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--topic bakwite-topic01 \
--partitions 3 \
--replication-factor 3
Created topic bakwite-topic01
创建后,查看topic列表
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--list
bakwite-topic01
· 查看topic
1,查看topic列表
| 参数 | 说明 |
|---|---|
| --list | 表是对topic的操作为:显示列表 |
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--list
bakwite-topic01
2,查看topic详情
| 参数 | 说明 |
|---|---|
| --describe | 表示要查看topic的详情 |
| --topic | 指定要查看的topic分区;如果不写,会查看所有 |
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--describe \
--topic bakwite-topic01 \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf
Topic: bakwite-topic01 TopicId: FGAY7jg1SPqLBXXMVpkuvw PartitionCount: 3 ReplicationFactor: 3 Configs: min.insync.replicas=2,unclean.leader.election.enable=false
Topic: bakwite-topic01 Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
Topic: bakwite-topic01 Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
Topic: bakwite-topic01 Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2
显示结果说明:
| 显示结果 | 二级字段 | 解释说明 |
|---|---|---|
| Topic: bakwite-topic01 | 表示topic名称 | |
| TopicId: FGAY7...... | 表示topic的ID | |
| PartitionCount: 3 | 表示该topic的分区数 | |
| ReplicationFactor: 3 | 表示该topic的副本数 | |
| Configs: | min.insync.replicas=2, | 最小同步副本数 生产者写入数据时,必须两个副本写入完成,才算成功 |
| unclean.leader.election.enable=false | 当leader宕机时,是否允许从非同步副本中选举新leader | |
| leader | 当前分区的leader所在broker节点; | |
| Isr | 当前处于同步状态的副本;因该是1.2.3都有才算完整; 说明,备份没有问题; |
3,各节点查看目录
[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0 #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:43 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2
[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:43 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1 #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:36 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2
[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:36 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2 # 主
· 修改topic
| 参数 | 说明 |
|---|---|
| --alter | 修改topic主题 |
注意:
- 只能修改分区数,副本数一经创建,无法修改;
- 修改分区,也只能增加,不能减少;
# 将分区数从:3,修改为:5;
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--alter \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--topic bakwite-topic01 \
--partitions 5
修改后查看topic详情
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--describe \
--topic bakwite-topic01 \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf
Topic: bakwite-topic01 TopicId: PN4Ew3jXRWmloK-sRUA80Q PartitionCount: 5 ReplicationFactor: 3 Configs: min.insync.replicas=2,unclean.leader.election.enable=false
Topic: bakwite-topic01 Partition: 0 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
Topic: bakwite-topic01 Partition: 1 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2
Topic: bakwite-topic01 Partition: 2 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
Topic: bakwite-topic01 Partition: 3 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
Topic: bakwite-topic01 Partition: 4 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2
查看各个节点目录验证
[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:43 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2 #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-3
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-4
[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0 #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:43 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:36 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-3 #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-4
[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1 #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:36 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-3
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-4 #主
· 删除topic
| 参数 | 说明 |
|---|---|
| --delete | 表示对topic进行删除操作 |
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--delete \
--topic bakwite-topic01 \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf
2,生产者写入数据
· 编辑生产者配置文件
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/producer.properties
# 链接kafaka的入口
bootstrap.servers=kafka01:9092,kafka02:9092,kafka03:9092
# 消息压缩类型【none】表示不压缩
# 可以设置为:gzip,snappy、lz4、zstd,但是会增加CPU开销
# 推荐【snappy】或【lz4】
compression.type=snappy
# 消息从发送到收到确认的总超时时间(2分钟)
delivery.timeout.ms=120000
# 单次请求(发送或获取元数据)的超时时间(30秒)
request.timeout.ms=30000
# 启用幂等生产者;
enable.idempotence=true
# # 开启幂等性,默认包含下面两项参数;
# # 消息确认机制,保证数据不丢失的关键配置
# # all表示所有ISR副本确认,才算写入成功,否则写入失败;
# acks=all
# # 当确定发送失败后的重试次数【下面这个值表示无限重试】
# retries=2147483647
# 批次和缓冲区
# 最多等待5ms打包
linger.ms=5
# 16KB批次大小
batch.size=16384
# 32MB发送缓冲区
buffer.memory=33554432
# 单个连接上允许的未确认请求最大数量
max.in.flight.requests.per.connection=5
# 指定key和value的序列化器
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer
# 链接集群用户
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="123456";
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
· 创建topic
[root@kafka01 ~ ]# cd /tools/kafka_2.13-3.3.2/
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--create \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--topic bakwite-topic01 \
--partitions 3 \
--replication-factor 3
Created topic bakwite-topic01.
· 生产者写入数据
1,交互式写入数据
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
--bootstrap-server 10.0.0.111:9092 \
--producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
--topic bakwite-topic01
>hello,world!
>bakwite!
>go on kfka!
>
2,非交互式写入数据
# 单行数据写入
[root@kafka01 ~ ]# echo "hello,bakwite!" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
--bootstrap-server 10.0.0.111:9092 \
--producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
--topic bakwite-topic01
# 多行写入
[root@kafka01 ~ ]# cat <<EOF | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
--bootstrap-server 10.0.0.111:9092 \
--producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
--topic bakwite-topic01
hello,wa
kafka new message
gogo!
EOF
# 从文件中读取写入
[root@kafka01 ~ ]# cat 1.txt
hei,kafka
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
--bootstrap-server 10.0.0.111:9092 \
--producer.config login-kafka.conf \
--topic bakwite-topic01 < 1.txt
3,消费者读取数据
· 编辑消费者配置文件
1,编辑消费者配置文件
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/consumer.properties
# 连接
bootstrap.servers=10.0.0.111:9092,10.0.0.112:9092,10.0.0.113:9092
# 消费者组ID
group.id=wa-group
# 可靠性
# 是否自动提交offset(消费位移)
enable.auto.commit=true
# 当没有初始offset或offset失效时,从哪里开始消费
# latest:表示最新的消息开始消费
# earliest:表示最早的消息开始消费(重头开始)
# none:没有offset时抛出异常
auto.offset.reset=earliest
# 性能
# 消费者从broker拉取数据的最小字节数
# 设置大了可以减少请求,但是会增加延迟;
fetch.min.bytes=1
# 等待fetch.min.bytes数据的最长时间(毫秒)
fetch.max.wait.ms=500
# 每个partition一次拉取的最大字节数(默认1MB)
max.partition.fetch.bytes=1048576
# 消费者一次从kafka拉取消息的最大数量;
max.poll.records=1000
# 容错
# 消费者会话超时时间(30秒)
session.timeout.ms=30000
# 消费者带一次拉取完毕后,超过5分钟,kafka会认为这个消费者挂了,触发冲平衡,重新分配消费者;
max.poll.interval.ms=300000
# 消费者定期发送心跳证明自己存活
heartbeat.interval.ms=3000
# 反序列化
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
# 链接集群用户
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="123456";
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
2,分发其他节点
[root@kafka01 ~ ]# scp /tools/kafka_2.13-3.3.2/config/consumer.properties 10.0.0.112:/tools/kafka_2.13-3.3.2/config/consumer.properties
[root@kafka01 ~ ]# scp /tools/kafka_2.13-3.3.2/config/consumer.properties 10.0.0.113:/tools/kafka_2.13-3.3.2/config/consumer.properties
· 读取所有数据
| 参数 | 说明 |
|---|---|
| --from-beginning | 从头开始读取数据(不写这个参数,默认值读新的写入数据) |
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
--bootstrap-server 10.0.0.111:9092 \
--consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
--topic bakwite-topic01 \
--from-beginning
hello,world!
bakwite!
go on kfka!
hello,bakwite!
hello,wa
kafka new message
gogo!
hei,kafka
· 读取前5条数据
| 参数 | 说明 |
|---|---|
| --max-messages 5 | 读取多少条数据 |
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
--bootstrap-server 10.0.0.111:9092 \
--consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
--topic bakwite-topic01 \
--from-beginning \
--max-messages 5
hello,world!
bakwite!
go on kfka!
hello,bakwite!
hello,wa
Processed a total of 5 messages
· 消费者组读取数据
1,生产者写入数据原理
【生产者】写入数据:
- 只写入topic主题的主分区中
- 写入是否成功,取决于topic的配置;
- 副本分区至少成功复制完多少个,算作写入成功?配置文件参数:【min.insync.replicas=2】;
- 也就是说,加上主分区,共计3个副本数据,成功写入2个,才算作数据写入完成;

2,消费者读取数据原理
【消费者】读取数据:
消费者读取数据,只读取【leader主分区副本数据】
- kafka为了保证数据的强一致性,就是leader分区与follower分区数据相同;
- 要求读取写入都必须是leader的分区;
- 因为,在读取时,副本follower很可能还没复制完,导致数据不一致,读取不全;
消费者读数据,会在kafka的data目录下创建读取记录的目录;
- 记录着读取数据的offset偏移量;
- 这些offset偏移量目录,可以让集群监控到各个节点的消费及情况;
- 即便集群宕机重启,消费者根据这些offset记录也能知道从哪里继续消费读数据;
- 避免数据重复读取;
[root@kafka01 ~ ]# ls -ld /tools/kafka_2.13-3.3.2/data/__consumer_offsets-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-10
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-11
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-12
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-13
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-14
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-15
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-16
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-17
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-18
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-19
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-2
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-20
......

3,kafka消费者读数据瓶颈
| 消费者读数据瓶颈问题 | 说明 |
|---|---|
| broker单点故障 | 当消费者在一个broker节点正在读数据,突然这台机器挂了,消费者也就挂了,整个消费停止; |
| 高并发数据积压 | 当topic有大量数据积压时,一个消费者消费性能有限,读取缓慢; |
4,消费者组
为了解决【消费者单点故障】与【高并发数据积压】的问题:
- kafak采用了比较特别的手段:【消费者组】;
- 就是创建消费者组:启动多个消费者,绑定这个消费者组;
- kafka会将topic内的所有分区,自动分配给组内消费者成员进行消费;
- 如果有新的消费者加入消费者组,会自动触发重平衡,再次平均分配分区;
- 组内【消费者】与【分区】是【1对多】的关系;
- 一个分区最多有一个消费者;
- 而一个消费者最多可以消费多个分区;
- 当触发重平衡后,组内成员会自动接管上一任消费者的offset,保正数据不会重复消费;
- 也就是说:
- 当我启动三个消费者进程分别在三台broker主机上;
- 就可以解决单点故障问题,和高并发数据积压问题;

将用户文件分发给其他节点
[root@kafka01 ~ ]# scp login-kafka.conf 10.0.0.112:~/
[root@kafka01 ~ ]# scp login-kafka.conf 10.0.0.113:~/
创建topic
# 创建topic
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--create \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--topic bakwite-topic01 \
--partitions 3 \
--replication-factor 3
Created topic bakwite-topic01
三个节点分别启动消费者,并【指定相同消费者组】
- 注意:消费者组会自动创建
- 当消费者组未被使用时,超过7天,自动删除,无需人为干预;
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
--bootstrap-server 10.0.0.111:9092 \
--consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
--topic bakwite-topic01
[root@kafka02 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
--bootstrap-server 10.0.0.111:9092 \
--consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
--topic bakwite-topic01
[root@kafka03 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
--bootstrap-server 10.0.0.111:9092 \
--consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
--topic bakwite-topic01
查看消费者组列表
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-consumer-groups.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--all-groups \
--list
wa-group
查看消费者组详细状态

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-consumer-groups.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--group wa-group \
--describe
生产者,多次写入数据
- 总结:
- 生产者每次写入数据,组内只有一个消费者会读取,不会重复消费;
- 即便一台机器挂了,也不影响消费数据;
[root@kafka01 ~ ]# echo "这是一条测试消息" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
--bootstrap-server 10.0.0.111:9092 \
--producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
--topic bakwite-topic01
[root@kafka01 ~ ]# echo "这是一条测试消息" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
--bootstrap-server 10.0.0.111:9092 \
--producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
--topic bakwite-topic01
[root@kafka01 ~ ]# echo "这是一条测试消息" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
--bootstrap-server 10.0.0.111:9092 \
--producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
--topic bakwite-topic01
[root@kafka01 ~ ]# echo "这是一条测试消息" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
--bootstrap-server 10.0.0.111:9092 \
--producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
--topic bakwite-topic01

4,kafka数据相关
· 什么情况下数据会丢
当生产者向leader分区写入数据后:
- 副本分区开始同步leader主分区的数据;
- 当一个broker节点尚未同步完成,leader节点突然宕机,同时kafka将这个主节点选举为leader分区;
- 此时,所有其他节点开始与这个broker节点同步数据(因为它是leader了),导致数据不完整;

【ISR】:主题分区同步情况;
- 当我们查看主题详情,看到最后最后一列叫做【ISR】
- 【In-Sync Replicas】 = 同步副本
- ISR这一列,就可以看到,当前副本分区的同步情况;
# 查看主题topic详情
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
--describe \
--topic bakwite-topic01 \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf
Topic: bakwite-topic01 TopicId: ZAz2sFLxTnaXr6GtRtErGw PartitionCount: 3 ReplicationFactor: 3 Configs: min.insync.replicas=2,segment.bytes=1073741824,retention.ms=604800000,unclean.leader.election.enable=false
Topic: bakwite-topic01 Partition: 0 Leader: 3 Replicas: 3,1,2 Isr: 3,2,1
Topic: bakwite-topic01 Partition: 1 Leader: 1 Replicas: 1,2,3 Isr: 3,2,1
Topic: bakwite-topic01 Partition: 2 Leader: 2 Replicas: 2,3,1 Isr: 3,2,1
· kafka如何保证数据不丢失
1,生产者配置文件
目标让生产者写入数据时:
启动幂等性:【enable.idempotence=true】
必须所有副本确认同步完成后,才算写入成功;
否则,判定失败,生产者重新写入数据;
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/producer.properties
......
# 消息从发送到收到确认的总超时时间(2分钟)
delivery.timeout.ms=120000
# 单次请求(发送或获取元数据)的超时时间(30秒)
request.timeout.ms=30000
# 启用幂等生产者;
enable.idempotence=true
# # 开启幂等性,默认包含下面两项参数;
# # 消息确认机制,保证数据不丢失的关键配置
# # all表示所有ISR副本确认,才算写入成功,否则写入失败;
# acks=all
# # 当确定发送失败后的重试次数【下面这个值表示无限重试】
# retries=2147483647
......
2,kafka配置文件
目标让kafka保证主分区宕机时:
- 必须同步完成的副本分区才可以被选举为主分区;
- 只有同步完成的副本分区,才有资格成为leader分区;
- 【unclean.leader.election.enable=false】
- 保证有同步完的备选分区
- 生产者写入时,必须起码有2个副本分区同步完成,才算做写入完成;
- 【min.insync.replicas=2】
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
...
# 必须两个副本同步完成,才算写入成功,否则,写入失败(回复生产者写入失败);
min.insync.replicas=2
# topic副本设置为3;保证数据安全;
default.replication.factor=3
# 主分区宕机时,不允许未同步完成同步的副本分区成为leader分区
unclean.leader.election.enable=false
# 是否自动平衡leader分布(消费者重平衡)
auto.leader.rebalance.enable=true
# 副本滞后时间阈值:30秒(超过则认为副本同步失效)
replica.lag.time.max.ms=30000
· kafka的kraft选举机制
1,Contrller与Broker角色
Controller节点核心功能:
- 主题与分区管理:负责增删改查“主题topic”、分区扩容工作;
- 分区Leader选举:当broker宕机时,Controller会为受到影响的分区快速选举出新的Leader副本,确保服务可用;
- 集群节点扩缩容:Controller能够监听Broker的加入和退出,并执行响应的“善后”工作;
- 元数据管理:Controller维护着最全的集群元数据,并负责将这些信息同步给集群内所有的Broker;
| kafka节点类型 | 名称 | 作用 |
|---|---|---|
| broker | 集群数据存储节点 | kafka的集群节点 |
| Controller | 控制节点 | kafka集群的管理功能节点 |
2,Zookeeper与Kraft集群模式区别
Zookeeper模式下集群模式:
- Broker在Zookeeper集群上创建临时节点来竞争Controller;
- 由Zookeeper的ZAB协议保证只有一个可以选举成功;
- 元数据(ISR、分区、主题)存储位置:
- 在Zookeeper中,Contrller需要频繁读取Zookeeper获取元数据;
- 再同步给其他Broker节点;
Kraft模式下集群模式:
- Controller不再是1个节点,而是一组(3、5个);
- 由Raft协议投片选举出Contrller,其他参与选举的节点作为备用,
- 主节点宕机,备用节点快速接管;
- 元数据(ISR、分区、主题)存储位置:
- 存储在kafka的一个元数据主题中:__cluster_metadata
- 由raft协议复制到所有Contrller节点;
- 本质上就是一个内置的、高可用的元数据日志(如MySQL的binlog),不再需要Zookeeper;
| 集群模式 | 对应 |
|---|---|
| Kraft | kafka3.3.2之后版本 |
| Zookeeper | kafka3.3.2之前版本,kafka4.0移除Zookeeper; |
5,kfka环境要求
| 场景 | broker | kraft-Controller | 副本 | 副本确认 | 吞吐量 | CPU/内存 | 存储容量 |
|---|---|---|---|---|---|---|---|
| 小型 | 3 | 3(与broker合并) | 3 | 2 | 100M/s | 4核/8G | <1T |
| 中小型 | 5-10 | 3(独立) | 3 | 2 | 100-500M/s | 8核/64G | 1-20T |
| 大型 | 10-20 | 3/5(独立) | 3 | 2 | 500-2000M/s | 32核/64G | 20TB-1PB |
| 超大 | 2-50 | 5(独立) | 3 | 2 | >2000M/s | 32核/64G | PB级别 |
五、ELFK架构实战
1,filebeat配置写入kafka
· 编辑filebeat配置文件
[root@kibana ~ ]# vim /tools/filebeat-8.18.8-linux-x86_64/filebeat.yml
filebeat.inputs:
- type: filestream
id: bakwite_01
enabled: true
paths:
- /test/1.txt
# 输出到kafka
output.kafka:
hosts:
- 10.0.0.111:9092
- 10.0.0.112:9092
- 10.0.0.113:9092
# 写入到topic主题
topic: bakwite-topic01
#【幂等性】
# 保证副本同步完成才算写入成功,否则重新写入;
required_acks: -1
# 只要失败就无限次重新写入,保证数据不丢失
max_retries: -1
# 写入超过30s算失败,重写
timeout: 30s
# 压缩方式
compression: gzip
# 轮询kafka节点写入
partition.round_robin:
# true表示:只写入存活的broker节点;
reachable_only: true
# ===== SSL =====
ssl.enabled: true
ssl.certificate_authorities:
- /tools/filebeat-8.18.8-linux-x86_64/certs/kafka-ca.crt
# ===== SASL =====
sasl.mechanism: PLAIN
username: filebeat
password: "123456"
· 复制kafka的服务证书
[root@kibana ~ ]# scp 10.0.0.111:/tools/kafka_2.13-3.3.2/certs/kafka-ca.crt /tools/filebeat-8.18.8-linux-x86_64/certs/
· 重启filebeat
[root@kibana ~ ]# systemctl restart filebeat.service
· kafka启动消费者
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
--bootstrap-server 10.0.0.111:9092 \
--consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
--topic bakwite-topic01
· 输入文件写入内容
[root@kibana ~ ]# echo "hello filebeat" >> /test/1.txt
[root@kibana ~ ]# echo "hello filebeat" >> /test/1.txt
[root@kibana ~ ]# echo "hello filebeat" >> /test/1.txt
· 查看消费者终端
{
"@timestamp":"2025-12-15T05:25:32.977Z",
"@metadata":{"beat":"filebeat","type":"_doc","version":"8.18.8"},
"host":{"name":"kibana"},
"agent":{
"ephemeral_id":"7c76994b-bc88-4f30-9e11-322dc4bfd151",
"id":"2f97e00a-7e4c-45d4-8bd1-a65742eaf4f8",
"name":"kibana",
"type":"filebeat",
"version":"8.18.8"
},
"log":{
"file":{
"path":"/test/1.txt",
"device_id":"64768",
"inode":"4587523"
},
"offset":1128
},
"message":"hello filebeat",
"input":{"type":"filestream"},
"ecs":{"version":"8.0.0"}
}

2,logstash读取kafka数据
· 生成JVM信任证书
1,复制kafka的ca证书
[root@kibana ~ ]# scp 10.0.0.111:/tools/kafka_2.13-3.3.2/certs/kafka-ca.crt /tools/logstash-8.18.8/certs/
2,使用JAVA证书工具生成证书
[root@kibana ~ ]# /tools/jdk1.8.0_461/bin/keytool -import \
-alias kafka-ca \
-file /tools/logstash-8.18.8/certs/kafka-ca.crt \
-keystore /tools/logstash-8.18.8/certs/kafka.truststore.jks \
-storepass 123456
.......
Trust this certificate? [no]: y # 这里需要确认【y】
#查看生成的证书
[root@kibana ~ ]# ls -l /tools/logstash-8.18.8/certs/
......
-rw-r--r-- 1 root root 1506 Dec 15 05:46 kafka.truststore.jks
· 编辑logstash配置文件
[root@kibana ~ ]# vim /tools/logstash-8.18.8/conf.d/test.conf
input {
kafka {
bootstrap_servers => "10.0.0.111:9092,10.0.0.112:9092,10.0.0.113:9092"
topics => ["bakwite-topic01"]
group_id => "xjzw-gp"
# ===== 安全 =====
security_protocol => "SASL_SSL"
sasl_mechanism => "PLAIN"
# 用户名密码
sasl_jaas_config => "org.apache.kafka.common.security.plain.PlainLoginModule required username='logstash' password='123456';"
# JVM专有证书
ssl_truststore_location => "/tools/logstash-8.18.8/certs/kafka.truststore.jks"
ssl_truststore_password => "123456"
# 从最新的offset开始消费
# earliest 表示从头开始消费
# latest 表示只读最新的,之前的不读;
auto_offset_reset => "earliest"
# 自动提交消费offset信息
enable_auto_commit => true
# 消费者线程(线程越多速度越快,跟cpu核心有关)
consumer_threads => 3
# 最慢3秒拉取一次消息
poll_timeout_ms => 3000
# 一次最多拉取500条消息
max_poll_records => 500
}
}
filter {
mutate {
remove_field => ["log","host","event","@version","host"]
}
}
output {
elasticsearch {
hosts => ["https://10.0.0.101:10200", "https://10.0.0.102:10200", "https://10.0.0.103:10200"]
index => "kafka-%{+YYYY.MM.dd}"
ssl => true
ssl_certificate_verification => true
cacert => "/tools/logstash-8.18.8/certs/ca.crt"
user => "elastic"
password => "123456"
}
}
· 重启logstash
[root@kibana ~ ]# systemctl restart logstash.service
· 写入数据测试
1,filebeat输入文件写入内容
[root@kibana ~ ]# echo "hello bakwite" >> /test/1.txt
2,查看kibana的索引

3,查看消费者组验证
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-consumer-groups.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--all-groups \
--list
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-consumer-groups.sh \
--bootstrap-server 10.0.0.111:9092 \
--command-config login-kafka.conf \
--group xjzw-gp \
--describe
END
Discussion
评论