Zookeeper优化与Kafka集群实战

7115 字
36 分钟
Zookeeper优化与Kafka集群实战
Zookeeper优化与Kafka集群实战

Zookeeper优化 & Kafka集群实战#

[TOC]


环境规划#

在上一篇笔记中我们已经在 Elk01/02/03 上部署好了 ZK 集群(10.0.0.6/7/8),本篇继续:

  • 第一步:对 ZK 集群做 JVM 调优 + 部署 zkUI 图形化
  • 第二步:在 ZK 之上部署 Kafka 集群,完成从部署到脚本、原理、优化、图形化的全套实战
主机IP角色
Elk0110.0.0.6ES + Kibana + ZK + Kafka(broker.id=1) + zkUI + Kafbat UI
Elk0210.0.0.7ES + ZK + Kafka(broker.id=2)
Elk0310.0.0.8ES + ZK + Kafka(broker.id=3)
Warning

🤡 实验前先关机拍快照,实验完关机恢复快照,确保环境干净 🤡

Terminal window
jiuzhao@Ubuntu ~$ virsh shutdown Elk01 Elk02 Elk03
jiuzhao@Ubuntu ~$ for v in Elk01 Elk02 Elk03; do virsh snapshot-create-as $v Elk-before-kafka; done
# Elk01 内存小,扩到 4G
jiuzhao@Ubuntu ~$ virsh setmaxmem Elk01 --size 4G --config
jiuzhao@Ubuntu ~$ virsh dumpxml Elk01 | grep -i memory
<memory unit='KiB'>4194304</memory>
<currentMemory unit='KiB'>2097152</currentMemory>
jiuzhao@Ubuntu ~$ virsh setmem Elk01 --size 4G --config
# 关机状态下,加 --config 选项为下一次启动预设内存大小
# 开机状态下,不需要加 --config 选项,立即生效
jiuzhao@Ubuntu ~$ virsh dumpxml Elk01 | grep -i memory
<memory unit='KiB'>4194304</memory>
<currentMemory unit='KiB'>4194304</currentMemory>
jiuzhao@Ubuntu ~$ virsh start Elk01 Elk02 Elk03
Important

⚠️ 2G 内存是最大约束——三台机都跑着 ES 集群(每台 512m 堆),ZK 默认堆高达 1000m 所以必须先做 ZK JVM 调优,腾出内存才能上 Kafka


Zookeeper 的 JVM 调优#

查看默认堆内存#

Terminal window
1)查看当前 ZK 的堆内存
root@Elk01 ~# ps -ef | grep -i zookeeper | grep Xmx
root ... java ... -Xmx1000m ... org.apache.zookeeper.server.quorum.QuorumPeerMain ...
# ZK 3.8 默认堆是 1000m(1G),对 2G 的机器太奢侈了
# 🔥 此刻没有 zkCli 客户端,所以只有一个进程 → 主类 QuorumPeerMain = 服务端
2)启动一个 zkCli 客户端,进程变成两个
root@Elk01 ~# zkCli.sh -server 10.0.0.6:2181
# 再开一个终端查看:
root@Elk01 ~# ps -ef | grep -i zookeeper | grep Xmx
... -Xmx1000m ... org.apache.zookeeper.server.quorum.QuorumPeerMain ← 服务端 1000m
... -Xmx256m ... org.apache.zookeeper.ZooKeeperMain ← 客户端 256m
# 多出来的主类 ZooKeeperMain = zkCli 客户端
'ps 命令很长,不用全看,拖到最后一行看主类名就行 ✅️'
3)用 jps 看更清爽
root@Elk01 ~# jps
613 Elasticsearch # ES 的 Java 进程
1535 QuorumPeerMain # ZK 服务端
1800 ZooKeeperMain # ZK 客户端
1911 Jps # jps 自己
# jps = Java Virtual Machine Process Status Tool(JDK 自带)
# 列出所有 Java 进程:PID + 主类名,一眼分清谁是干嘛的
Tip

💡 jps -l 能显示完整主类名;多 JDK 环境 jps 可能不在 PATH(先 source /etc/profile.d/zk.sh 对应的环境变量文件,或写全路径 /usr/share/elasticsearch/jdk/bin/jps)

Note

📌 ZK 有两类 Java 进程,所以下面调优时两个堆要一起调(服务端占大头,客户端小头):

主类角色堆来源默认值
QuorumPeerMain==ZK 服务端==ZK_SERVER_HEAP1000m
ZooKeeperMain==ZK 客户端==(zkCli)ZK_CLIENT_HEAP256m

修改堆内存#

Terminal window
1)修改 zkEnv.sh(服务器 + 客户端堆)
root@Elk01 ~# vim /usr/local/apache-zookeeper-3.8.6-bin/bin/zkEnv.sh
...
# ZK_SERVER_HEAP="${ZK_SERVER_HEAP:-1000}" # 服务器堆默认 1000
ZK_SERVER_HEAP="${ZK_SERVER_HEAP:-256}" # 改成 256
# ZK_CLIENT_HEAP="${ZK_CLIENT_HEAP:-256}" # 客户端堆默认 256
ZK_CLIENT_HEAP="${ZK_CLIENT_HEAP:-128}" # 改成 128
...
2)同步配置文件到集群其他节点
root@Elk01 ~# scp /usr/local/apache-zookeeper-3.8.6-bin/bin/zkEnv.sh 10.0.0.7:/usr/local/apache-zookeeper-3.8.6-bin/bin/
root@Elk01 ~# scp /usr/local/apache-zookeeper-3.8.6-bin/bin/zkEnv.sh 10.0.0.8:/usr/local/apache-zookeeper-3.8.6-bin/bin/
3)滚动重启所有节点(一台一台来,保证 quorum 不丢)
root@Elk01 ~# source /etc/profile.d/zk.sh && zkServer.sh stop && zkServer.sh start
root@Elk02 ~# source /etc/profile.d/zk.sh && zkServer.sh stop && zkServer.sh start
root@Elk03 ~# source /etc/profile.d/zk.sh && zkServer.sh stop && zkServer.sh start
4)再次检查 JVM 是否生效(客户端也开着,两个堆一起验证)
root@Elk01 ~# ps -ef | grep -i zookeeper | grep Xmx
... -Xmx256m ... org.apache.zookeeper.server.quorum.QuorumPeerMain ← 服务端 256m ✅
... -Xmx128m ... org.apache.zookeeper.ZooKeeperMain ← 客户端 128m ✅
root@Elk01 ~# jps -l
2321 jdk.jcmd/sun.tools.jps.Jps
2164 org.apache.zookeeper.server.quorum.QuorumPeerMain # 服务端 1000→256m ✅
2261 org.apache.zookeeper.ZooKeeperMain # 客户端 256→128m ✅
613 org.elasticsearch.bootstrap.Elasticsearch
Tip

💡 生产环境 ZK 堆建议设置 2GB+,4GB 即可;这里实验环境机器太小(2G),才压到 256m

字段原值调优后说明
ZK_SERVER_HEAP1000256服务器堆,占内存大头
ZK_CLIENT_HEAP256128客户端工具堆,够用就行

zkUI — Zookeeper 图形化管理#

Note

老师原本讲的是 zkWeb(该项目已停止维护,下载源失效),这里用更主流的 zkUI(DeemOpen/zkui)替代 zkUI 是一个 Java Web 应用,默认端口 9090,支持 ZK 节点的 CRUD 操作(增、删、改、查) C (Create):创建、R (Read):查询/读取、U (Update):更新/修改、D (Delete):删除

部署步骤#

Terminal window
1)编译 zkUI(源码需要 Maven 编译)
jiuzhao@Ubuntu ~$ git clone https://github.com/DeemOpen/zkui.git
jiuzhao@Ubuntu ~$ mvn --version
找不到命令 “mvn”
jiuzhao@Ubuntu zkui$ source /etc/profile.d/maven.sh
jiuzhao@Ubuntu ~$ cd zkui && MAVEN_OPTS="--add-opens java.base/java.net=ALL-UNNAMED" mvn clean package
....................
[INFO] BUILD SUCCESS
[INFO] --------------------------------------
[INFO] Total time: 4.366 s
[INFO] Finished at: 2026-08-04T14:21:20+08:00
[INFO] --------------------------------------
# 产物:target/zkui-2.0-SNAPSHOT-jar-with-dependencies.jar
jiuzhao@Ubuntu zkui$ mv ./target/{zkui-2.0-SNAPSHOT-jar-with-dependencies,zkui}.jar
jiuzhao@Ubuntu zkui$ scp config.cfg ./target/zkui.jar Elk01:/tmp/
# config.cfg 是 zkUI 的配置文件(核心内容就是指定 zkUI 连哪个 ZK 集群)
2)zkUI 需要一个 JDK8 运行(机器上 ES 的是 JDK22,不保险)
jiuzhao@Ubuntu ~$ wget -O /home/jiuzhao/下载/jdk8-temurin.tar.gz https://api.adoptium.net/v3/binary/latest/8/ga/linux/x64/jdk/hotspot/normal/eclipse
# 下载 Temurin JDK8 到宿主机,再 scp 到 Elk01
jiuzhao@Ubuntu ~$ scp /home/jiuzhao/下载/jdk8-temurin.tar.gz Elk01:/tmp
3)Elk01 上解压 JDK8 + 部署 zkui
root@Elk01 ~# tar xf /tmp/jdk8-temurin.tar.gz -C /usr/local/
root@Elk01 ~# mkdir -p /soft/zkui && cp /tmp/zkui.jar /tmp/config.cfg /soft/zkui/
4)修改 config.cfg 指向 ZK 集群
root@Elk01 ~# vim /soft/zkui/config.cfg
serverPort=9090 # Web 端口
zkServer=10.0.0.6:2181,10.0.0.7:2181,10.0.0.8:2181 # ZK 集群地址(逗号分隔)
userSet = {"users": [{"username":"admin","password":"manager","role":"ADMIN"}]}
# 默认账号 admin/manager
5)启动 zkUI(JDK8 运行)
root@Elk01 ~# cd /soft/zkui
# zkUI 在当前工作目录找 config.cfg;所以先 cd /soft/zkui
root@Elk01 zkui# nohup /usr/local/jdk8u502-b07/bin/java -Xmx256m -jar /soft/zkui/zkui.jar &>> /tmp/zkui.log &
root@Elk01 zkui# ss -lnt | grep 9090
LISTEN 0 50 *:9090 *:* # 监听成功
6)访问 WebUI
http://10.0.0.6:9090 # 账号 admin / 密码 manager

zkUI 界面
zkUI 界面

Tip

💡 浏览器打开 zkUI 后,左侧能看到 ZK 的节点树,支持增删改查节点 后续部署 Kafka 后,可以在这里直接看到 Kafka 在 ZK 里生成的元数据节点

Terminal window
# zkUI 这个管理页面也挺拉的!
jiuzhao@Ubuntu ~$ scp /home/jiuzhao/下载/zkWeb-v1.2.1.jar Elk01:/tmp
root@Elk01 ~# mkdir -p /soft/zkWeb && cp /tmp/zkWeb-v1.2.1.jar /soft/zkWeb/
root@Elk01 ~# nohup /usr/local/jdk8u502-b07/bin/java -Xmx256m -jar /soft/zkWeb/zkWeb-v1.2.1.jar &>> /tmp/zkWeb.log &
root@Elk01 ~# tail -1 /tmp/zkWeb.log
[2026-08-04 15:47:56 INFO main StartupInfoLogger.java:59] c.y.zkweb.ZkWebSpringBootApplication --> Started ZkWebSpringBoot...
root@Elk01 ~# ss -ntl | grep 8099
LISTEN 0 100 *:8099 *:*
root@Elk01 ~# jps
706 QuorumPeerMain
611 Elasticsearch
1944 zkui.jar
2047 zkWeb-v1.2.1.jar
# 浏览器访问
http://10.0.0.6:8099

添加Zookeeper集群
添加Zookeeper集群

Tip

描述信息、集群IP<2181>(逗号分隔)、延迟时间(3000ms == 3s)

zkWeb 界面
zkWeb 界面


Kafka 消息队列#

MQ 是什么#

Note

通俗解释:消息队列(Message Queue)就是用来缓存数据的中间件,多用于高并发、数据量大的场景

  • 常见 MQ:ActiveMQ、RocketMQ、RabbitMQ、Kafka 等

消息队列三大优势:

  • 削峰填谷:突发流量先进队列,系统按自己的节奏慢慢消化
  • 异步提速:发完消息就返回,不用等对方处理完
  • 架构解耦:生产者和消费者互不依赖,各自升级互不影响

🌰 把 Kafka 理解成生活里的菜鸟驿站——你和快递员是两个程序,驿站就是 MQ 能替代驿站的还有蜂巢、快递站(RabbitMQ、RocketMQ…),Kafka 只是其中最能打的那个

生产者

Producer

消息队列 MQ

Kafka

生产者

Producer

生产者

Producer

消费者

Consumer

消费者

Consumer

生产者

Producer

消息队列 MQ

Kafka

生产者

Producer

生产者

Producer

消费者

Consumer

消费者

Consumer

Important

📌 版本演进(重要)

  • Kafka 2.8+:首次引入了 KRaft 模式,但这只是一个预览(Preview)版本
  • Kafka 3.x:ZK 模式和 KRaft 模式并存,生产主流还是 ZK 模式(经典)
  • Kafka 4.0+:完全不依赖 Zookeeper,只支持 KRaft

🎯 为什么先学 3.x + ZK?——Kafka+ZK 是最经典的组合,也最有难度 4.x 把 ZK 收编成一套系统,简单了,但底层原理(选举、元数据)都藏在 KRaft 里了 先把经典的 ZK 模式吃透,再去看 4.x 就一通百通

kafka下载
kafka下载

单点部署#

Terminal window
1)下载 kafka 软件包(3.x 系列最新稳定版 3.9.2)
jiuzhao@Ubuntu ~$ wget -P /home/jiuzhao/下载 https://archive.apache.org/dist/kafka/3.9.2/kafka_2.13-3.9.2.tgz
2)拷贝并解压
jiuzhao@Ubuntu ~$ scp /home/jiuzhao/下载/kafka_2.13-3.9.2.tgz Elk01:/tmp
root@Elk01 ~# tar xf /tmp/kafka_2.13-3.9.2.tgz -C /usr/local/
root@Elk01 ~# ls /usr/local/kafka_2.13-3.9.2/
bin config libs LICENSE ...
3)配置环境变量
root@Elk01 ~# vim /etc/profile.d/kafka.sh
#!/bin/bash
export KAFKA_HOME=/usr/local/kafka_2.13-3.9.2
export JAVA_HOME=/usr/share/elasticsearch/jdk # 复用 ES 自带 JDK22
export PATH=$PATH:$KAFKA_HOME/bin:$JAVA_HOME/bin
root@Elk01 ~# source /etc/profile.d/kafka.sh
4)预调 JVM 堆(2G 机器避免 OOM)
root@Elk01 ~# sed -i 's/KAFKA_HEAP_OPTS="-Xmx1G -Xms1G"/KAFKA_HEAP_OPTS="-Xmx256m -Xms256m"/' \
$KAFKA_HOME/bin/kafka-server-start.sh
# 老师讲义是先 1G 启动、后面再优化
# 但我们的机器 ES+ZK 已占 ~1.3G,1G 堆启动必 OOM,所以把优化前置了
5)修改配置文件 server.properties
root@Elk01 ~# vim $KAFKA_HOME/config/server.properties
...
broker.id=1 # 唯一标识 broker 节点
advertised.listeners=PLAINTEXT://10.0.0.6:9092 # 当前节点监听地址+端口
# PLAINTEXT 明文传输,也没有认证
log.dirs=/var/lib/kafka # 指定 kafka 数据目录
zookeeper.connect=10.0.0.6:2181,10.0.0.7:2181,10.0.0.8:2181/kafka-v3.9.2 # ZK集群地址
root@Elk01 ~# mkdir -p /var/lib/kafka
6)启动 kafka 单点
root@Elk01 ~# kafka-server-start.sh -daemon $KAFKA_HOME/config/server.properties
-daemon # 后台启动
root@Elk01 ~# tail /usr/local/kafka_2.13-3.9.2/logs/server.log
...
INFO [KafkaServer id=1] started (kafka.server.KafkaServer)
...
root@Elk01 ~# ss -lnt | grep 9092
LISTEN 0 50 *:9092 *:* # 9092 是 kafka 的对外端口
配置项值说明
broker.id1==每个 broker 集群内唯一标识==,集群内不能重复
advertised.listenersPLAINTEXT://IP:9092告诉客户端”来这个地址找我”,必须写真实 IP
log.dirs/var/lib/kafkakafka 数据落盘目录(顺序写磁盘)
zookeeper.connect三台:2181/kafka-v3.9.2元数据存 ZK,/kafka-v3.9.2 是==chroot 隔离路径==
Tip

💡 zookeeper.connect 里带 /kafka-v3.9.2 后缀 = 把 kafka 的元数据隔离在 ZK 的独立子路径下 这样多个 Kafka 集群可以共用一套 ZK,互不干扰 kafka 启动时会在 ZK 自动创建这个路径,无需手动建

/kafka-v3.9.2
/kafka-v3.9.2

集群部署#

Terminal window
1)拷贝程序到其他节点 + 同步环境变量
root@Elk01 ~# scp -r /usr/local/kafka_2.13-3.9.2/ 10.0.0.7:/usr/local/
root@Elk01 ~# scp -r /usr/local/kafka_2.13-3.9.2/ 10.0.0.8:/usr/local/
root@Elk01 ~# scp /etc/profile.d/kafka.sh 10.0.0.7:/etc/profile.d/
root@Elk01 ~# scp /etc/profile.d/kafka.sh 10.0.0.8:/etc/profile.d/
2)其他节点修改配置文件
root@Elk02 ~# source /etc/profile.d/kafka.sh
root@Elk02 ~# vim ${KAFKA_HOME}/config/server.properties
broker.id=2
advertised.listeners=PLAINTEXT://10.0.0.7:9092
root@Elk03 ~# source /etc/profile.d/kafka.sh ; vim ${KAFKA_HOME}/config/server.properties
broker.id=3
advertised.listeners=PLAINTEXT://10.0.0.8:9092
3)启动 kafka 服务
root@Elk02 ~# kafka-server-start.sh -daemon $KAFKA_HOME/config/server.properties
root@Elk02 ~# ss -lnt | grep 9092
LISTEN 0 50 *:9092 *:*
root@Elk03 ~# kafka-server-start.sh -daemon $KAFKA_HOME/config/server.properties
root@Elk03 ~# ss -lnt | grep 9092
LISTEN 0 50 *:9092 *:*
4)zkWeb 验证 — kafka 在 ZK 里注册的 broker
root@Elk01 ~# zkCli.sh -server 10.0.0.6:2181 ls /kafka-v3.9.2/brokers/ids
[1, 2, 3] # 三个 broker 全部注册 ✅️
'在 zkWeb 的节点树里也能看到 /kafka-v3.9.2/brokers/ids/1|2|3'

ids
ids

Kafka Brokers

Zookeeper

Elk01:2181

Elk02:2181

Elk03:2181

Elk01 (id=1):9092

Elk02 (id=2):9092

Elk03 (id=3):9092

Kafka Brokers

Zookeeper

Elk01:2181

Elk02:2181

Elk03:2181

Elk01 (id=1):9092

Elk02 (id=2):9092

Elk03 (id=3):9092

Kafka systemd 托管#

手动 kafka-server-start.sh -daemon 启动的进程无开机自启,重启 VM 后要手动拉起 对标 ZK 的 systemd 服务,给 Kafka 也配一个 kafka.service,实现开机自启 + 优雅停止

Terminal window
1)编写 systemd 服务文件(三台一样,放在 /etc/systemd/system/)
root@Elk01 ~# cat > /etc/systemd/system/kafka.service <<"EOF"
[Unit]
# 服务描述
Description=Apache Kafka Server
# 官方文档地址
Documentation=https://kafka.apache.org/
# 等网络和 zookeeper 就绪后再启动
After=network-online.target zookeeper.service
Wants=network-online.target
[Service]
# kafka-server-start.sh 内部 fork 出 java 守护进程,用 forking
Type=forking
# 环境变量:KAFKA_HOME / JAVA_HOME(复用 ES 自带 JDK)
Environment=KAFKA_HOME=/usr/local/kafka_2.13-3.9.2
Environment=JAVA_HOME=/usr/share/elasticsearch/jdk
Environment=PATH=/usr/local/kafka_2.13-3.9.2/bin:/usr/share/elasticsearch/jdk/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin
# 启动命令
ExecStart=/usr/local/kafka_2.13-3.9.2/bin/kafka-server-start.sh -daemon /usr/local/kafka_2.13-3.9.2/config/server.properties
# 停止命令(优雅停止,等待数据同步完成,不要 kill -9)
ExecStop=/usr/local/kafka_2.13-3.9.2/bin/kafka-server-stop.sh
# 异常退出才重启
Restart=on-failure
RestartSec=3
[Install]
# 开机自启
WantedBy=multi-user.target
EOF
2)三台停手动进程 → 转 systemd 接管
root@Elk01 ~# kafka-server-stop.sh
root@Elk01 ~# systemctl daemon-reload && systemctl enable --now kafka
root@Elk01 ~# systemctl is-active kafka
active
root@Elk01 ~# ss -lnt | grep 9092
LISTEN 0 50 *:9092 *:*
✅️ Kafka 由 systemd 接管,开机自启
# Elk02 / Elk03 执行相同操作
3)验证集群完整性(systemd 管理下)
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --list
__consumer_offsets
kpyun-linux
root@Elk01 ~# zkCli.sh -server 10.0.0.6:2181 ls /kafka-v3.9.2/brokers/ids
[1, 2, 3] # 三个 broker 全部注册 ✅️
Tip

💡 为什么 ExecStop 用 kafka-server-stop.sh 而不是 systemctl stop 默认行为? kafka-server-stop.sh 会优雅停止:先等数据同步完成再关闭,避免丢数据 直接 kill -9 强杀 = 数据同步中断,可能丢消息(详见下文”Kafka 为何丢失数据”章节)

⚠️ Type=forking:kafka-server-start.sh 内部会 fork 出 java 守护进程再退出 systemd 用 forking 判断”主进程 fork 成功 = 服务已启动”,与 ZK 服务文件同款

生产者与消费者验证#

Terminal window
1)启动生产者 — topic kpyun-linux 不存在时会自动创建
# 需 broker 开启 auto.create.topics.enable=true,默认开启
root@Elk01 ~# kafka-console-producer.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-linux
>www.kpyun.com # 写入第 1 条消息
'首次创建 topic 时 broker 要先选 leader、分配副本,会出现 LEADER_NOT_AVAILABLE 警告'
# 自动重试后成功,属正常现象
>jzops # 写入第 2 条消息
>xixi # 写入第 3 条消息
>haha # 写入第 4 条消息
>
2)查看 topic 列表 — 验证 topic 已创建
root@Elk01 ~# source /etc/profile.d/kafka.sh
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --list
kpyun-linux
3)启动消费者 — 实时拉取新消息(pull 模型,消费者主动向 broker 请求数据)
root@Elk02 ~# kafka-console-consumer.sh --bootstrap-server 10.0.0.7:9092 --topic kpyun-linux
# 此时在生产者端继续输入 666 / 777 / 888,消费者会实时打印 ✅️
4)从头消费 — --from-beginning 重置 offset 到最早位置,拉取全部历史消息
root@Elk03 ~# kafka-console-consumer.sh --bootstrap-server 10.0.0.8:9092 --topic kpyun-linux --from-beginning
www.kpyun.com # offset=0
jzops # offset=1
xixi # offset=2
haha # offset=3
666 # offset=4(步骤 3 生产者新增)
777 # offset=5
888 # offset=6 — 7 条历史数据全部拿到 ✅️
5)再次查看 topic 列表 — __consumer_offsets 出现说明消费者组已注册
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --list
__consumer_offsets # Kafka 内置 topic,存储消费者组的消费位置
kpyun-linux
Warning

⚠️ __consumer_offsets 是 kafka 的内置 topic,专门存消费者组的 offset(消费位置) 属于系统内部数据,不要试图删除!

命令作用
kafka-console-producer.sh --bootstrap-server ... --topic xxx启动生产者(交互式输入)
--bootstrap-server指定 Kafka broker 地址,多节点逗号分隔即可,不用全写
kafka-console-consumer.sh --bootstrap-server ... --topic xxx启动消费者(实时消费)
kafka-console-consumer.sh ... --from-beginning从头消费全部历史数据
kafka-topics.sh --bootstrap-server ... --list列出集群中所有 topic

Kafka 常用术语#

kafka常用术语架构图
kafka常用术语架构图

术语含义
brokerkafka 集群的每一个节点
kafka cluster也叫 “broker list”,整个 kafka 集群
producer==生产者==,向 kafka 集群写入数据的一方
consumer==消费者==,向 kafka 集群读取数据的一方
topic==主题==,逻辑概念,producer/consumer 数据读写的逻辑分类单元
partition==分区==,逻辑概念,topic 之下进一步切分的并行单元,最少 1 个
offset==偏移量==,分区里每条消息的位置编号,消费者靠它记录消费进度
replica==副本==,物理概念,真正负责磁盘存储和主从同步
Important

📌 逻辑概念 vs 物理存储

  • topic 和 partition 是逻辑概念,只是消息的分类方式和并行度划分,不直接存数据
  • replica(副本)才是真正存数据的东西 —— 每条消息最终落到某个副本的磁盘上
  • 可以这样理解:topic → 文件夹名(逻辑),partition → 子文件夹名(逻辑)
    • replica → 子文件夹里实际的 .log / .index 文件(物理)
Important

📌 副本(replica)的 leader/follower 关系

  • 每个分区都有 1 个 leader 副本 + N 个 follower 副本
  • leader:对外承担 ==读写请求==,是唯一可读写的那份数据
  • follower:只管从 leader ==同步数据==,不对外服务,leader 挂了才顶上
Note
  • ES 可以设 replica = 0(无副本),数据只靠主分片单份存储,挂了就丢
  • Kafka 必须至少 1 个分区 + 1 个副本,设计上不允许无副本
    • 虽然 replication-factor=1 只有一个副本即 leader,但实质上也是 replica
  • ES 副本可分担读请求;Kafka follower 纯备份,不参与读写
broker.id节点IP
1Elk0110.0.0.6
2Elk0210.0.0.7
3Elk0310.0.0.8
Important

看输出时对照这个表就行! broker.id **只能是整数 **✅,没办法用字符串命名 ❌

topic 管理脚本#

Terminal window
1)查看 topic 列表
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --list
__consumer_offsets
kpyun-linux
2)查看指定 topic 详细信息
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-linux --describe
Topic: kpyun-linux PartitionCount: 1 ReplicationFactor: 1
Topic: kpyun-linux Partition: 0 Leader: 2 Replicas: 2 Isr: 2
PartitionCount: # 该 topic 有多少个分区
ReplicationFactor: # 每个分区有多少个副本(1 表示只有 leader,无 follower)
# 1 个分区 1 个副本,leader 在 broker 2 上
3)创建 topic
3.1 默认创建(1 分区 1 副本)
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic xixi --create
Created topic xixi.
3.2 创建指定分区数(3 分区)
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic haha --create --partitions 3
Created topic haha.
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic haha --describe
Topic: haha Partition: 0 Leader: 3 Replicas: 3
Topic: haha Partition: 1 Leader: 1 Replicas: 1
Topic: haha Partition: 2 Leader: 2 Replicas: 2
# 3 个分区自动分散到 3 台 broker(Leader: 3/1/2),这就是==分区的分布式存储==
3.3 创建指定分区和副本数量(5 分区 2 副本)
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic hehe --create --partitions 5 --replication-factor 2
Created topic hehe.
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic hehe --describe
Topic: hehe PartitionCount: 5 ReplicationFactor: 2
Topic: hehe Partition: 0 Leader: 2 Replicas: 2,1 Isr: 2,1
Topic: hehe Partition: 1 Leader: 3 Replicas: 3,2 Isr: 3,2
Topic: hehe Partition: 2 Leader: 1 Replicas: 1,3 Isr: 1,3
Topic: hehe Partition: 3 Leader: 2 Replicas: 2,3 Isr: 2,3
Topic: hehe Partition: 4 Leader: 3 Replicas: 3,1 Isr: 3,1
# 每个分区 2 个副本,leader 和 follower 交叉分布,任何一个 broker 挂了都不丢数据
`broker.id 为 1 的机器是 Elk01 `
# Elk01 在分区 0、2、4 上(看副本哪些有Elk01的broker.id)
root@Elk01 ~# ll -d /var/lib/kafka/hehe-*
drwxr-xr-x 2 root root 4096 Aug 5 09:12 /var/lib/kafka/hehe-0/
drwxr-xr-x 2 root root 4096 Aug 5 09:12 /var/lib/kafka/hehe-2/
drwxr-xr-x 2 root root 4096 Aug 5 09:12 /var/lib/kafka/hehe-4/
# hehe-0 就是 topic=hehe、分区 0 的物理存储目录,leader 或 follower 都存一份
# 目录名 {topic}-{partitionId} 是 Kafka 的固定命名规则
Tip

💡 --describe 输出解读:

  • Leader:当前谁负责读写这个分区
  • Replicas:该分区副本分布在哪几个 broker
  • Isr:与 leader 数据同步(数据保持一致)的副本集合(详见后面”丢数据原理”)
Terminal window
4)修改分区数量(仅支持调大,不支持减小!)
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic hehe --alter --partitions 10
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic hehe --alter --partitions 3
Error while executing topic command : Topic currently has 10 partitions, which is higher than the requested 3.
# ❌️ 分区只能增不能减 —— 减少分区会导致已有数据无法重新分布
'分区的数据分布已经写死在磁盘了,缩分区 = 要挪数据,Kafka 不允许'
5)删除 topic
5.1 删除单个
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic hehe --delete
5.2 删除多个(逗号分割)
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic haha,xixi --delete
操作命令关键参数说明
查看列表--list列出所有 topic
查看详情--topic xxx --describe分区数、副本、leader、Isr
创建--create可选 --partitions / --replication-factor
改分区--alter --partitions N==只能调大,不能调小==
删除--delete支持逗号分割一次删多个

消费者组#

kafka生产者和消费者架构图解
kafka生产者和消费者架构图解

Tip

📌 图解说明(结合上图阅读)

1. Topic 与 Partition 的物理关系

  • oldboyedu-linux 是 ==topic(逻辑名称)==,生产者往这个主题写,消费者从这个主题读
  • -0、-1、-2 分别是 topic 下的 ==3 个 partition==,数据按 key hash 分散落入不同分区
  • 每个分区内消息有独立的 ==offset==(位置编号)

2. __consumer_offsets — 消费进度的”书签”

  • 消费者每读完一条消息,Kafka 会把 offset 记入内置 topic: __consumer_offsets
  • 图中 broker.id=92 下方的小云朵就是 offset 记录:消费者组 linux99 在分区 oldboyedu-linux-1 上已经消费到 offset 12
  • 作用:哪怕消费者挂了(c3 宕机),消费者组 linux99 中另一个成员(c2)可以以同一组身份从 offset 的下一条接着读,不丢不重

3. 消费者组的角色隔离

  • 图中两组消费者互不干扰:linux100 组(c1)和 linux99 组(c2/c3)各自独立维护自己的 offset
  • 同一个消费者组内,partition 只能分配给一个消费者,但不同消费者组可以同时消费同一个 partition
    • 📌 注意:一个消费者实例(JVM 进程)只能有一个 group.id,不能同时属于两个消费者组
  • 消费者组成员变化(有人挂了 / 有人加入)→ 触发 rebalance,partition 重新分配,拿到新分区的消费者仍从该分区的 offset 继续读
Important

📌 消费者组(consumer group)三条铁律

  1. 同一个消费者组的消费者,==不能同时消费同一个 topic 的同一个 partition==
  2. 消费者组成员增减时,触发 rebalance(重新分配 partition)
  3. partition 数量变化时,同样触发 rebalance

topic: kpyun-k8s(5 分区)

消费者组 group=jzops

消费者 C1

(加入)

消费者 C2

(加入)

P0

P1

P2

P3

P4

topic: kpyun-k8s(5 分区)

消费者组 group=jzops

消费者 C1

(加入)

消费者 C2

(加入)

P0

P1

P2

P3

P4

消费者组脚本实战#

Terminal window
1)准备 topic(5 分区 2 副本)
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-k8s --create --partitions 5 --replication-factor 2
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-k8s --describe
Topic: kpyun-k8s PartitionCount: 5 ReplicationFactor: 2 Configs:
Topic: kpyun-k8s Partition: 0 Leader: 1 Replicas: 1,2 Isr: 1,2
Topic: kpyun-k8s Partition: 1 Leader: 2 Replicas: 2,3 Isr: 2,3
Topic: kpyun-k8s Partition: 2 Leader: 3 Replicas: 3,1 Isr: 3,1
Topic: kpyun-k8s Partition: 3 Leader: 1 Replicas: 1,3 Isr: 1,3
Topic: kpyun-k8s Partition: 4 Leader: 2 Replicas: 2,1 Isr: 2,1
2)启动生产者写入数据
root@Elk01 ~# kafka-console-producer.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-k8s
>111111111111111111111
>22222222222222222222222
>3333333333333333333
>
3)启动消费者并指定消费者组,从头消费
root@Elk01 ~# kafka-console-consumer.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-k8s --from-beginning --group jzops
111111111111111111111
22222222222222222222222
3333333333333333333
4)查看消费者组列表
root@Elk01 ~# kafka-consumer-groups.sh --bootstrap-server 10.0.0.6:9092 --list
jzops
5)查看消费者组详情(一个消费者拿了所有分区)
root@Elk01 ~# kafka-consumer-groups.sh --bootstrap-server 10.0.0.6:9092 --group jzops --describe ;echo
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
jzops kpyun-k8s 0 0 0 0 console-consumer-55fe336f... /10.0.0.6 console-consumer
jzops kpyun-k8s 1 3 3 0 console-consumer-55fe336f... /10.0.0.6 console-consumer
jzops kpyun-k8s 2 0 0 0 console-consumer-55fe336f... /10.0.0.6 console-consumer
jzops kpyun-k8s 3 0 0 0 console-consumer-55fe336f... /10.0.0.6 console-consumer
jzops kpyun-k8s 4 0 0 0 console-consumer-55fe336f... /10.0.0.6 console-consumer
# 3 条消息全进了 partition 1(CURRENT-OFFSET=3, LAG=0)
# 一个消费者实例拿全部 5 个分区 → CONSUMER-ID 全是同一个,也在同一台机器上 /10.0.0.6
# 消费者退出后 → CONSUMER-ID、HOST、CLIENT-ID 全变 -,表示无活跃成员在线
字段含义
GROUP消费者组名称(--group 指定的那个)
TOPIC被消费的 topic
PARTITION该消费者组分到的分区编号
CURRENT-OFFSET==当前消费到的位置==(该分区已经读到了第几条)
LOG-END-OFFSET==该分区最新消息的末尾位置==(总共写到了第几条)
LAG==积压量== = LOG-END-OFFSET − CURRENT-OFFSET,值越大说明消费者越跟不上
CONSUMER-ID消费者实例 ID,- 表示当前无活跃消费者
HOST消费者所在机器 IP,- 表示无活跃消费者
CLIENT-ID客户端标识(连接时传入的 client.id),- 表示无活跃消费者
Terminal window
`退出所有消费者后`
6)再次写入数据,观察 LAG 累积
root@Elk01 ~# echo -e "aaaaaaaaa\nb\nc\nd\ne\nf" | kafka-console-producer.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-k8s
root@Elk01 ~# kafka-consumer-groups.sh --bootstrap-server 10.0.0.6:9092 --group jzops --describe
jzops kpyun-k8s 2 0 6 6 - - -
# 新数据进了 partition 2,LOG-END-OFFSET=6,但没消费者 → LAG=6(积压)
# LAG = LOG-END-OFFSET - CURRENT-OFFSET,也就是"还没被消费的消息数"
Caution

⚠️ LAG 是 kafka 运维最重要的指标——LAG 持续变大 = 消费者处理不过来了,0 才是健康的 排查方向:消费者挂了 / 消费者进程太少 / 消费逻辑太慢

Terminal window
7)启动第二个消费者(同组)→ 触发 rebalance
root@Elk01 ~# kafka-console-consumer.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-k8s --from-beginning --group jzops # 后台跑
root@Elk01 ~# kafka-consumer-groups.sh --bootstrap-server 10.0.0.6:9092 --group jzops --describe
jzops kpyun-k8s 0 0 0 0 console-consumer-634e31cf... /10.0.0.6 console-consumer
jzops kpyun-k8s 1 3 3 0 console-consumer-634e31cf... /10.0.0.6 console-consumer
jzops kpyun-k8s 2 6 6 0 console-consumer-634e31cf... /10.0.0.6 console-consumer
jzops kpyun-k8s 3 0 0 0 console-consumer-6e509e8a... /10.0.0.6 console-consumer
jzops kpyun-k8s 4 0 0 0 console-consumer-6e509e8a... /10.0.0.6 console-consumer
# 观察 CONSUMER-ID:出现两个不同的消费者
# 消费者1 拿了 0/1/2 分区,消费者2 拿了 3/4 分区 → 这就是 rebalance ✅️
'第二个消费者一加入,partition 自动重新分配,一个分区只会被组内一个消费者消费'

Kafka 为何丢失数据#

kafka为何丢失数据图解
kafka为何丢失数据图解

📌 先记五个术语

术语全称含义
ISRIn-Sync Replicas==和 leader 数据完全同步==的副本集合
OSROut-of-Sync Replicas和 leader 数据不同步(落后)的副本集合
ARAll Replica s所有副本 = leader + follower,AR = ISR + OSR
LEOLog End Offset某个分区副本的的最后一条消息的 offset
HWHigh Watermark==ISR 中最小的 LEO==,消费者只能消费 HW 之前的数据
Important

📌 HW 的真正含义 — 为什么是最小的 LEO?

  • HW 取当前 ISR 列表中最小的 LEO,消费者只能消费 HW 之前的数据
  • ISR 的检测是周期性的(默认 replica.lag.time.max.ms = 30s),不是实时的
  • 关键场景:在两次检测之间,某个 ISR 副本可能已经跟不上 leader 的同步速度了,但仍然留在 ISR 列表中
  • 此时这个”掉队但还在 ISR”的副本,它的 LEO 比 leader 小 → HW 就以它的 LEO 为准
  • 等到下一个 30s 检测周期:如果它还没追上 → 被踢出 ISR,进入 OSR;追上了 → 继续留在 ISR
  • 因此 HW 可以理解为:在周期检测的间隙,以 ISR 内最慢的那个副本为准,保证消费者看到的数据在所有 ISR 副本上都是完整的
Note

📌 丢数据的本质 — 什么情况下会丢?

丢数据只有一种情况:leader 挂了,新 leader 还没同步完旧 leader 的最新数据

  1. 生产者写入消息到 leader → leader 返回 ACK 给生产者
  2. 在 follower 还没来得及同步这条消息时,leader 突然挂了(kill -9、断电)
  3. ISR 中某个 follower 被选为新 leader
  4. 但这个新 leader 没拿到旧 leader 刚写的那条消息 → 消息丢失
  5. 生产者以为写入成功了,实际新 leader 没有这条数据

⚠️ 为什么很难丢?

  • 只要不是 kill -9 强行杀死或服务器突然断电,几乎不会丢
  • 正常的 kafka-server-stop.sh 停止会等待数据同步完成再关闭
  • 大数据领域某个特性是价值密度低,少量消息丢失对业务影响可控

避免方案:

  • acks=all(-1)— 等所有 ISR 副本都确认写入后才算成功
  • min.insync.replicas >= 2 — ISR 至少保持 2 个副本,防止只有 leader 一个副本时挂了
  • 停止 Kafka 用 kafka-server-stop.sh,不要 kill -9

📌 关于 acks=all 的补充说明

acks 值含义可靠性
0生产者不等待确认,发出去就行最低,消息可能没到 broker 就丢了
1(默认)等 leader 确认写入磁盘即可中等,leader 挂了且没同步给 follower 则丢
all / -1等所有 ISR 副本都确认写完最高,只要 ISR 有一个副本活着就不会丢
  • acks 是生产者端的配置
    • 命令行指定:kafka-console-producer.sh ... --producer-property acks=all
  • acks=all 必须配合 min.insync.replicas 一起用才有意义:
    • 比如 replication-factor=3, min.insync.replicas=2
      • 3 份数据副本,ISR 至少 2 个副本在线,才允许写入
      • 如果 ISR 只剩 leader 一个(另两个挂了),acks=all 时写入直接失败,宁可不写也不丢数据
    • 副本数 ≥ min.insync.replicas,且 min.insync.replicas ≥ 2,这样才能真正防丢

Kafka 优化篇#

JVM 调优#

Terminal window
1)查看当前 JVM 大小
root@Elk01 ~# ps -ef | grep kafka | grep Xms
... java -Xmx256m -Xms256m ...
# 部署时已前置调优到 256m(2G 机器)
2)修改配置文件(生产环境建议 6GB)
root@Elk01 ~# vim `which kafka-server-start.sh`
if [ "x$KAFKA_HEAP_OPTS" = "x" ]; then
# export KAFKA_HEAP_OPTS="-Xmx1G -Xms1G"
export KAFKA_HEAP_OPTS="-Xmx256m -Xms256m"
fi
3)同步到其他节点 + 重启
root@Elk01 ~# scp `which kafka-server-start.sh` 10.0.0.7:/usr/local/kafka_2.13-3.9.2/bin/
root@Elk01 ~# scp `which kafka-server-start.sh` 10.0.0.8:/usr/local/kafka_2.13-3.9.2/bin/
4)验证
root@Elk01 ~# ps -ef | grep kafka | grep Xms
... -Xmx256m -Xms256m ... ✅️
Tip

💡 生产环境 Kafka JVM 建议 6GB 即可,数据是顺序写磁盘的,堆不用开太大 磁盘 IO 建议 1.5w 转速机械盘起步,条件允许直接上固态

禁用自动创建 topic#

Terminal window
1)修改所有节点的配置文件
root@Elk01 ~# echo 'auto.create.topics.enable=false' >> $KAFKA_HOME/config/server.properties
# 所有节点都执行 ⚠️
root@Elk01 ~# tail -1 $KAFKA_HOME/config/server.properties
auto.create.topics.enable=false
root@Elk03 ~# tail -1 $KAFKA_HOME/config/server.properties
auto.create.topics.enable=false
2)重启 kafka 集群
root@Elk01 ~# kafka-server-stop.sh ; sleep 5 ; kafka-server-start.sh -daemon $KAFKA_HOME/config/server.properties
# 所有节点都执行 ⚠️
3)测试验证(无法自动创建 topic)
root@Elk01 ~# kafka-console-producer.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-xixi
>1111
WARN The metadata response reported a recoverable issue ... {kpyun-xixi=UNKNOWN_TOPIC_OR_PARTITION}
# ❌️ UNKNOWN_TOPIC_OR_PARTITION = topic 不存在,禁止自动创建 ✅️
root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --list
__consumer_offsets
kpyun-k8s
kpyun-linux
# topic 没有被偷偷创建出来 ✅️
'生产环境必须关掉,防止开发手滑打错 topic 名,自动建出一堆垃圾 topic'

数据保留时间#

Terminal window
1)查看默认保留时间
root@Elk01 ~# grep log.retention.hours $KAFKA_HOME/config/server.properties
log.retention.hours=168
# 默认 168h = 7 天
Caution

⚠️ 不指定 log.retention.hours 默认保留 168h(7 天) 生产环境要和开发人员确认数据需要存多久,提前算好存储容量 检查 log.dirs 存储路径的磁盘空间是否够用

其他优化参数#

官方文档
官方文档

参考官方文档:Kafka 3.9.x 版本配置

核心思路:

  • CPU:多核心(压缩/解压、网络处理都吃 CPU)
  • 磁盘:顺序写 + 大容量 + 固态更好
  • 网络:带宽决定吞吐上限
  • 32GB 内存的服务器一般就够用了,kafka 不会吃太多内存

Kafbat UI — Kafka 图形化管理(Docker)#

Note

老师原本讲的是 kafka-eagle (EFAK)(2022 年后已停更,对 Kafka 3.9 兼容性存疑) 这里改用社区更活跃的 Kafbat UI(Provectus kafka-ui 的社区维护版,2026 年仍持续更新) Kafbat UI 是集中式 Web 应用,==单实例==通过 bootstrap-servers 连整个集群,不需要每台都装

部署步骤(Docker)#

Terminal window
1)Elk01 拉取镜像(只有 Elk01 需要装 Docker,Kafbat UI 单实例即可)
root@Elk01 ~# docker pull docker.xuanyuan.run/kafbat/kafka-ui:latest; docker tag docker.xuanyuan.run/kafbat/kafka-ui:latest kafbat/kafka-ui:latest; docker rmi docker.xuanyuan.run/kafbat/kafka-ui:latest
2)运行容器 — 关键是 KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS 指向三台 broker
root@Elk01 ~# docker run -d --name kafbat-ui --restart unless-stopped -p 18080:8080 \
-e KAFKA_CLUSTERS_0_NAME=my-elastic-cluster \
-e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=10.0.0.6:9092,10.0.0.7:9092,10.0.0.8:9092 \
-e DYNAMIC_CONFIG_ENABLED=true \
kafbat/kafka-ui:latest
'ZK 3.5+ 有个默认的 Jetty 管理端口就是 8080,容器映射时撞上了,换个宿主端口即可'
3)验证容器
root@Elk01 ~# docker ps | grep kafbat-ui
29c171ed6a40 kafbat/kafka-ui:latest Up ... 0.0.0.0:18080->8080/tcp kafbat-ui
4)访问 WebUI
http://10.0.0.6:18080/

Kafbat UI 界面
Kafbat UI 界面

  • 在 UI 里能看到 broker 列表、topic 详情、消费者组 offset、消息内容——比命令行直观得多
Tip

💡 Kafbat UI 的 API 也能验证集群:

Terminal window
root@Elk01 ~# curl -s http://10.0.0.6:18080/api/clusters | jq ;echo
[
{
"name": "my-elastic-cluster",
"defaultCluster": null,
"status": "ONLINE",
"lastError": null,
"brokerCount": 3,
"onlinePartitionCount": 61,
"topicCount": 4,
"bytesInPerSec": null,
"bytesOutPerSec": null,
"readOnly": false,
"version": "3.9-IV0",
"features": [
"TOPIC_DELETION",
"CLIENT_QUOTA_MANAGEMENT",
"FTS_ENABLED"
],
"controller": "ZOOKEEPER"
}
]
# brokerCount=3, controller=ZOOKEEPER → 成功连上 ZK 模式的 Kafka 3.9 集群 ✅️

文章分享

如果这篇文章对你有帮助,欢迎分享给更多人!

Profile Image of the Author
久棹
不是先学好了再干,而是先干起来再学习,干中学!
分类
站点统计
文章
115
分类
15
标签
272
总字数
306,562
运行时长
0 天
最后活动
0 天前
文章目录