Zookeeper优化与Kafka集群实战

Zookeeper优化 & Kafka集群实战
[TOC]
环境规划
在上一篇笔记中我们已经在 Elk01/02/03 上部署好了 ZK 集群(10.0.0.6/7/8),本篇继续:
- 第一步:对 ZK 集群做 JVM 调优 + 部署 zkUI 图形化
- 第二步:在 ZK 之上部署 Kafka 集群,完成从部署到脚本、原理、优化、图形化的全套实战
| 主机 | IP | 角色 |
|---|---|---|
| Elk01 | 10.0.0.6 | ES + Kibana + ZK + Kafka(broker.id=1) + zkUI + Kafbat UI |
| Elk02 | 10.0.0.7 | ES + ZK + Kafka(broker.id=2) |
| Elk03 | 10.0.0.8 | ES + ZK + Kafka(broker.id=3) |
🤡 实验前先关机拍快照,实验完关机恢复快照,确保环境干净 🤡
jiuzhao@Ubuntu ~$ virsh shutdown Elk01 Elk02 Elk03jiuzhao@Ubuntu ~$ for v in Elk01 Elk02 Elk03; do virsh snapshot-create-as $v Elk-before-kafka; done# Elk01 内存小,扩到 4Gjiuzhao@Ubuntu ~$ virsh setmaxmem Elk01 --size 4G --configjiuzhao@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⚠️ 2G 内存是最大约束——三台机都跑着 ES 集群(每台 512m 堆),ZK 默认堆高达 1000m 所以必须先做 ZK JVM 调优,腾出内存才能上 Kafka
Zookeeper 的 JVM 调优
查看默认堆内存
1)查看当前 ZK 的堆内存root@Elk01 ~# ps -ef | grep -i zookeeper | grep Xmxroot ... 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 ~# jps613 Elasticsearch # ES 的 Java 进程1535 QuorumPeerMain # ZK 服务端1800 ZooKeeperMain # ZK 客户端1911 Jps # jps 自己# jps = Java Virtual Machine Process Status Tool(JDK 自带)# 列出所有 Java 进程:PID + 主类名,一眼分清谁是干嘛的💡 jps -l 能显示完整主类名;多 JDK 环境 jps 可能不在 PATH(先 source /etc/profile.d/zk.sh 对应的环境变量文件,或写全路径 /usr/share/elasticsearch/jdk/bin/jps)
📌 ZK 有两类 Java 进程,所以下面调优时两个堆要一起调(服务端占大头,客户端小头):
| 主类 | 角色 | 堆来源 | 默认值 |
|---|---|---|---|
QuorumPeerMain | ==ZK 服务端== | ZK_SERVER_HEAP | 1000m |
ZooKeeperMain | ==ZK 客户端==(zkCli) | ZK_CLIENT_HEAP | 256m |
修改堆内存
1)修改 zkEnv.sh(服务器 + 客户端堆)root@Elk01 ~# vim /usr/local/apache-zookeeper-3.8.6-bin/bin/zkEnv.sh...# ZK_SERVER_HEAP="${ZK_SERVER_HEAP:-1000}" # 服务器堆默认 1000ZK_SERVER_HEAP="${ZK_SERVER_HEAP:-256}" # 改成 256# ZK_CLIENT_HEAP="${ZK_CLIENT_HEAP:-256}" # 客户端堆默认 256ZK_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 startroot@Elk02 ~# source /etc/profile.d/zk.sh && zkServer.sh stop && zkServer.sh startroot@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 -l2321 jdk.jcmd/sun.tools.jps.Jps2164 org.apache.zookeeper.server.quorum.QuorumPeerMain # 服务端 1000→256m ✅2261 org.apache.zookeeper.ZooKeeperMain # 客户端 256→128m ✅613 org.elasticsearch.bootstrap.Elasticsearch💡 生产环境 ZK 堆建议设置 2GB+,4GB 即可;这里实验环境机器太小(2G),才压到 256m
| 字段 | 原值 | 调优后 | 说明 |
|---|---|---|---|
ZK_SERVER_HEAP | 1000 | 256 | 服务器堆,占内存大头 |
ZK_CLIENT_HEAP | 256 | 128 | 客户端工具堆,够用就行 |
zkUI — Zookeeper 图形化管理
老师原本讲的是 zkWeb(该项目已停止维护,下载源失效),这里用更主流的 zkUI(DeemOpen/zkui)替代 zkUI 是一个 Java Web 应用,默认端口 9090,支持 ZK 节点的 CRUD 操作(增、删、改、查) C (Create):创建、R (Read):查询/读取、U (Update):更新/修改、D (Delete):删除
部署步骤
1)编译 zkUI(源码需要 Maven 编译)jiuzhao@Ubuntu ~$ git clone https://github.com/DeemOpen/zkui.gitjiuzhao@Ubuntu ~$ mvn --version找不到命令 “mvn”jiuzhao@Ubuntu zkui$ source /etc/profile.d/maven.shjiuzhao@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.jarjiuzhao@Ubuntu zkui$ mv ./target/{zkui-2.0-SNAPSHOT-jar-with-dependencies,zkui}.jarjiuzhao@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 到 Elk01jiuzhao@Ubuntu ~$ scp /home/jiuzhao/下载/jdk8-temurin.tar.gz Elk01:/tmp
3)Elk01 上解压 JDK8 + 部署 zkuiroot@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.cfgserverPort=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/zkuiroot@Elk01 zkui# nohup /usr/local/jdk8u502-b07/bin/java -Xmx256m -jar /soft/zkui/zkui.jar &>> /tmp/zkui.log &root@Elk01 zkui# ss -lnt | grep 9090LISTEN 0 50 *:9090 *:* # 监听成功
6)访问 WebUIhttp://10.0.0.6:9090 # 账号 admin / 密码 manager
💡 浏览器打开 zkUI 后,左侧能看到 ZK 的节点树,支持增删改查节点 后续部署 Kafka 后,可以在这里直接看到 Kafka 在 ZK 里生成的元数据节点
# 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 8099LISTEN 0 100 *:8099 *:*root@Elk01 ~# jps706 QuorumPeerMain611 Elasticsearch1944 zkui.jar2047 zkWeb-v1.2.1.jar
# 浏览器访问http://10.0.0.6:8099
描述信息、集群IP<2181>2181>(逗号分隔)、延迟时间(3000ms == 3s)

Kafka 消息队列
MQ 是什么
通俗解释:消息队列(Message Queue)就是用来缓存数据的中间件,多用于高并发、数据量大的场景
- 常见 MQ:ActiveMQ、RocketMQ、RabbitMQ、Kafka 等
消息队列三大优势:
- 削峰填谷:突发流量先进队列,系统按自己的节奏慢慢消化
- 异步提速:发完消息就返回,不用等对方处理完
- 架构解耦:生产者和消费者互不依赖,各自升级互不影响
🌰 把 Kafka 理解成生活里的菜鸟驿站——你和快递员是两个程序,驿站就是 MQ 能替代驿站的还有蜂巢、快递站(RabbitMQ、RocketMQ…),Kafka 只是其中最能打的那个
📌 版本演进(重要)
- 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下载
单点部署
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:/tmproot@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/bashexport KAFKA_HOME=/usr/local/kafka_2.13-3.9.2export JAVA_HOME=/usr/share/elasticsearch/jdk # 复用 ES 自带 JDK22export PATH=$PATH:$KAFKA_HOME/bin:$JAVA_HOME/binroot@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.propertiesroot@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 9092LISTEN 0 50 *:9092 *:* # 9092 是 kafka 的对外端口| 配置项 | 值 | 说明 |
|---|---|---|
broker.id | 1 | ==每个 broker 集群内唯一标识==,集群内不能重复 |
advertised.listeners | PLAINTEXT://IP:9092 | 告诉客户端”来这个地址找我”,必须写真实 IP |
log.dirs | /var/lib/kafka | kafka 数据落盘目录(顺序写磁盘) |
zookeeper.connect | 三台:2181/kafka-v3.9.2 | 元数据存 ZK,/kafka-v3.9.2 是==chroot 隔离路径== |
💡 zookeeper.connect 里带 /kafka-v3.9.2 后缀 = 把 kafka 的元数据隔离在 ZK 的独立子路径下
这样多个 Kafka 集群可以共用一套 ZK,互不干扰
kafka 启动时会在 ZK 自动创建这个路径,无需手动建

集群部署
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.shroot@Elk02 ~# vim ${KAFKA_HOME}/config/server.propertiesbroker.id=2advertised.listeners=PLAINTEXT://10.0.0.7:9092
root@Elk03 ~# source /etc/profile.d/kafka.sh ; vim ${KAFKA_HOME}/config/server.propertiesbroker.id=3advertised.listeners=PLAINTEXT://10.0.0.8:9092
3)启动 kafka 服务root@Elk02 ~# kafka-server-start.sh -daemon $KAFKA_HOME/config/server.propertiesroot@Elk02 ~# ss -lnt | grep 9092LISTEN 0 50 *:9092 *:*root@Elk03 ~# kafka-server-start.sh -daemon $KAFKA_HOME/config/server.propertiesroot@Elk03 ~# ss -lnt | grep 9092LISTEN 0 50 *:9092 *:*
4)zkWeb 验证 — kafka 在 ZK 里注册的 brokerroot@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'
Kafka systemd 托管
手动 kafka-server-start.sh -daemon 启动的进程无开机自启,重启 VM 后要手动拉起
对标 ZK 的 systemd 服务,给 Kafka 也配一个 kafka.service,实现开机自启 + 优雅停止
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.serviceWants=network-online.target
[Service]# kafka-server-start.sh 内部 fork 出 java 守护进程,用 forkingType=forking# 环境变量:KAFKA_HOME / JAVA_HOME(复用 ES 自带 JDK)Environment=KAFKA_HOME=/usr/local/kafka_2.13-3.9.2Environment=JAVA_HOME=/usr/share/elasticsearch/jdkEnvironment=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-failureRestartSec=3
[Install]# 开机自启WantedBy=multi-user.targetEOF
2)三台停手动进程 → 转 systemd 接管root@Elk01 ~# kafka-server-stop.shroot@Elk01 ~# systemctl daemon-reload && systemctl enable --now kafkaroot@Elk01 ~# systemctl is-active kafkaactiveroot@Elk01 ~# ss -lnt | grep 9092LISTEN 0 50 *:9092 *:*✅️ Kafka 由 systemd 接管,开机自启# Elk02 / Elk03 执行相同操作
3)验证集群完整性(systemd 管理下)root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --list__consumer_offsetskpyun-linuxroot@Elk01 ~# zkCli.sh -server 10.0.0.6:2181 ls /kafka-v3.9.2/brokers/ids[1, 2, 3] # 三个 broker 全部注册 ✅️💡 为什么 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 服务文件同款
生产者与消费者验证
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.shroot@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --listkpyun-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-beginningwww.kpyun.com # offset=0jzops # offset=1xixi # offset=2haha # offset=3666 # offset=4(步骤 3 生产者新增)777 # offset=5888 # 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⚠️ __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 常用术语

| 术语 | 含义 |
|---|---|
| broker | kafka 集群的每一个节点 |
| kafka cluster | 也叫 “broker list”,整个 kafka 集群 |
| producer | ==生产者==,向 kafka 集群写入数据的一方 |
| consumer | ==消费者==,向 kafka 集群读取数据的一方 |
| topic | ==主题==,逻辑概念,producer/consumer 数据读写的逻辑分类单元 |
| partition | ==分区==,逻辑概念,topic 之下进一步切分的并行单元,最少 1 个 |
| offset | ==偏移量==,分区里每条消息的位置编号,消费者靠它记录消费进度 |
| replica | ==副本==,物理概念,真正负责磁盘存储和主从同步 |
📌 逻辑概念 vs 物理存储
- topic 和 partition 是逻辑概念,只是消息的分类方式和并行度划分,不直接存数据
- replica(副本)才是真正存数据的东西 —— 每条消息最终落到某个副本的磁盘上
- 可以这样理解:topic → 文件夹名(逻辑),partition → 子文件夹名(逻辑)
- replica → 子文件夹里实际的
.log/.index文件(物理)
- replica → 子文件夹里实际的
📌 副本(replica)的 leader/follower 关系
- 每个分区都有 1 个 leader 副本 + N 个 follower 副本
- leader:对外承担 ==读写请求==,是唯一可读写的那份数据
- follower:只管从 leader ==同步数据==,不对外服务,leader 挂了才顶上
- ES 可以设
replica = 0(无副本),数据只靠主分片单份存储,挂了就丢 - Kafka 必须至少 1 个分区 + 1 个副本,设计上不允许无副本
- 虽然
replication-factor=1只有一个副本即 leader,但实质上也是 replica
- 虽然
- ES 副本可分担读请求;Kafka follower 纯备份,不参与读写
| broker.id | 节点 | IP |
|---|---|---|
| 1 | Elk01 | 10.0.0.6 |
| 2 | Elk02 | 10.0.0.7 |
| 3 | Elk03 | 10.0.0.8 |
看输出时对照这个表就行! broker.id **只能是整数 **✅,没办法用字符串命名 ❌
topic 管理脚本
1)查看 topic 列表root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --list__consumer_offsetskpyun-linux
2)查看指定 topic 详细信息root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-linux --describeTopic: kpyun-linux PartitionCount: 1 ReplicationFactor: 1 Topic: kpyun-linux Partition: 0 Leader: 2 Replicas: 2 Isr: 2PartitionCount: # 该 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 --createCreated topic xixi.
3.2 创建指定分区数(3 分区)root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic haha --create --partitions 3Created 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 2Created topic hehe.root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic hehe --describeTopic: 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 的固定命名规则💡 --describe 输出解读:
Leader:当前谁负责读写这个分区Replicas:该分区副本分布在哪几个 brokerIsr:与 leader 数据同步(数据保持一致)的副本集合(详见后面”丢数据原理”)
4)修改分区数量(仅支持调大,不支持减小!)root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic hehe --alter --partitions 10root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic hehe --alter --partitions 3Error 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 | 支持逗号分割一次删多个 |
消费者组

📌 图解说明(结合上图阅读)
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,不能同时属于两个消费者组
- 📌 注意:一个消费者实例(JVM 进程)只能有一个
- 消费者组成员变化(有人挂了 / 有人加入)→ 触发 rebalance,partition 重新分配,拿到新分区的消费者仍从该分区的 offset 继续读
📌 消费者组(consumer group)三条铁律
- 同一个消费者组的消费者,==不能同时消费同一个 topic 的同一个 partition==
- 消费者组成员增减时,触发 rebalance(重新分配 partition)
- partition 数量变化时,同样触发 rebalance
消费者组脚本实战
1)准备 topic(5 分区 2 副本)root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-k8s --create --partitions 5 --replication-factor 2root@Elk01 ~# kafka-topics.sh --bootstrap-server 10.0.0.6:9092 --topic kpyun-k8s --describeTopic: 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 jzops111111111111111111111222222222222222222222223333333333333333333
4)查看消费者组列表root@Elk01 ~# kafka-consumer-groups.sh --bootstrap-server 10.0.0.6:9092 --listjzops
5)查看消费者组详情(一个消费者拿了所有分区)root@Elk01 ~# kafka-consumer-groups.sh --bootstrap-server 10.0.0.6:9092 --group jzops --describe ;echoGROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-IDjzops kpyun-k8s 0 0 0 0 console-consumer-55fe336f... /10.0.0.6 console-consumerjzops kpyun-k8s 1 3 3 0 console-consumer-55fe336f... /10.0.0.6 console-consumerjzops kpyun-k8s 2 0 0 0 console-consumer-55fe336f... /10.0.0.6 console-consumerjzops kpyun-k8s 3 0 0 0 console-consumer-55fe336f... /10.0.0.6 console-consumerjzops 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),- 表示无活跃消费者 |
`退出所有消费者后`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-k8sroot@Elk01 ~# kafka-consumer-groups.sh --bootstrap-server 10.0.0.6:9092 --group jzops --describejzops kpyun-k8s 2 0 6 6 - - -# 新数据进了 partition 2,LOG-END-OFFSET=6,但没消费者 → LAG=6(积压)# LAG = LOG-END-OFFSET - CURRENT-OFFSET,也就是"还没被消费的消息数"⚠️ LAG 是 kafka 运维最重要的指标——LAG 持续变大 = 消费者处理不过来了,0 才是健康的 排查方向:消费者挂了 / 消费者进程太少 / 消费逻辑太慢
7)启动第二个消费者(同组)→ 触发 rebalanceroot@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 --describejzops kpyun-k8s 0 0 0 0 console-consumer-634e31cf... /10.0.0.6 console-consumerjzops kpyun-k8s 1 3 3 0 console-consumer-634e31cf... /10.0.0.6 console-consumerjzops kpyun-k8s 2 6 6 0 console-consumer-634e31cf... /10.0.0.6 console-consumerjzops kpyun-k8s 3 0 0 0 console-consumer-6e509e8a... /10.0.0.6 console-consumerjzops 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 为何丢失数据

📌 先记五个术语
| 术语 | 全称 | 含义 |
|---|---|---|
ISR | In-Sync Replicas | ==和 leader 数据完全同步==的副本集合 |
OSR | Out-of-Sync Replicas | 和 leader 数据不同步(落后)的副本集合 |
AR | All Replica s | 所有副本 = leader + follower,AR = ISR + OSR |
LEO | Log End Offset | 某个分区副本的的最后一条消息的 offset |
HW | High Watermark | ==ISR 中最小的 LEO==,消费者只能消费 HW 之前的数据 |
📌 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 副本上都是完整的
📌 丢数据的本质 — 什么情况下会丢?
丢数据只有一种情况:leader 挂了,新 leader 还没同步完旧 leader 的最新数据
- 生产者写入消息到 leader → leader 返回 ACK 给生产者
- 在 follower 还没来得及同步这条消息时,leader 突然挂了(
kill -9、断电) - ISR 中某个 follower 被选为新 leader
- 但这个新 leader 没拿到旧 leader 刚写的那条消息 → 消息丢失
- 生产者以为写入成功了,实际新 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 调优
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 ... ✅️💡 生产环境 Kafka JVM 建议 6GB 即可,数据是顺序写磁盘的,堆不用开太大 磁盘 IO 建议 1.5w 转速机械盘起步,条件允许直接上固态
禁用自动创建 topic
1)修改所有节点的配置文件root@Elk01 ~# echo 'auto.create.topics.enable=false' >> $KAFKA_HOME/config/server.properties# 所有节点都执行 ⚠️root@Elk01 ~# tail -1 $KAFKA_HOME/config/server.propertiesauto.create.topics.enable=falseroot@Elk03 ~# tail -1 $KAFKA_HOME/config/server.propertiesauto.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>1111WARN 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_offsetskpyun-k8skpyun-linux# topic 没有被偷偷创建出来 ✅️'生产环境必须关掉,防止开发手滑打错 topic 名,自动建出一堆垃圾 topic'数据保留时间
1)查看默认保留时间root@Elk01 ~# grep log.retention.hours $KAFKA_HOME/config/server.propertieslog.retention.hours=168# 默认 168h = 7 天⚠️ 不指定 log.retention.hours 默认保留 168h(7 天)
生产环境要和开发人员确认数据需要存多久,提前算好存储容量
检查 log.dirs 存储路径的磁盘空间是否够用
其他优化参数

参考官方文档:Kafka 3.9.x 版本配置
核心思路:
- CPU:多核心(压缩/解压、网络处理都吃 CPU)
- 磁盘:顺序写 + 大容量 + 固态更好
- 网络:带宽决定吞吐上限
- 32GB 内存的服务器一般就够用了,kafka 不会吃太多内存
Kafbat UI — Kafka 图形化管理(Docker)
老师原本讲的是 kafka-eagle (EFAK)(2022 年后已停更,对 Kafka 3.9 兼容性存疑)
这里改用社区更活跃的 Kafbat UI(Provectus kafka-ui 的社区维护版,2026 年仍持续更新)
Kafbat UI 是集中式 Web 应用,==单实例==通过 bootstrap-servers 连整个集群,不需要每台都装
部署步骤(Docker)
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 指向三台 brokerroot@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-ui29c171ed6a40 kafbat/kafka-ui:latest Up ... 0.0.0.0:18080->8080/tcp kafbat-ui
4)访问 WebUIhttp://10.0.0.6:18080/
- 在 UI 里能看到 broker 列表、topic 详情、消费者组 offset、消息内容——比命令行直观得多
💡 Kafbat UI 的 API 也能验证集群:
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 集群 ✅️文章分享
如果这篇文章对你有帮助,欢迎分享给更多人!















