消息中间件之 Kafka
Kafka 背景
Kafka是由LinkedIn开发并开源的分布式消息系统,因其分布式及高吞吐率而被广泛使用,现已与Cloudera Hadoop,Apache Storm,Apache Spark集成。
Kafka名字灵感来自卡夫卡小说(《变形记》作者),寓意:数据像小说情节一样有序流动但错综复杂。

2010年,Kafka诞生。2011年,正式开源。2012年:成为Apache顶级项目。LinkedIn用Kafka解决了自己的数据洪流问题,却意外创造了改变互联网基础设施的工具。
2015,国内开始引进;2019,巅峰;2023年,AI崛起。
MQ 消息队列
MQ(Message Queue,消息队列)就像是一个“快递站”,负责把消息从一个系统传递到另一个系统,确保数据不丢失、不重复,还能应对高并发流量。
没有MQ:
用户点击 → 服务器(崩溃)
有MQ:
用户点击 → MQ(缓冲) → 每秒放100个请求给服务器
想象你是一个卖家(生产者),每天要发很多快递给买家(消费者)。
没有驿站(MQ):
你每次都要亲自联系快递员,等他上门取件。
如果快递员忙不过来,你的包裹就会堆积,甚至丢失。
有驿站(MQ):
你把包裹送到驿站,驿站帮你暂存、排队,快递员有空时来取。
好处:你不用等快递员,快递员也不用立刻接单,双方更高效。
因此,我们总结一下,关于MQ的核心功能。
表格 还在加载中,请等待加载完成后再尝试复制
点对点消息系统
在点对点系统中,消息被保留在队列中。 一个或多个消费者可以消耗队列中的消息,但是特定消息只能由最多一个消费者消费。 一旦消费者读取队列中的消息,它就从该队列中消失。

发布-订阅消息系统
在发布 - 订阅系统中,消息被保留在主题中。 与点对点系统不同,消费者可以订阅一个或多个主题并使用该主题中的所有消息。
在发布 - 订阅系统中,消息生产者称为发布者,消息使用者称为订阅者。
一个现实生活的例子是电视,它发布不同的渠道,如运动,电影,音乐等,任何人都可以订阅自己的频道集,并获得他们订阅的频道时可用。

Kafka 概述
Kafka 介绍
Kafka是一种分布式的,基于发布/订阅的消息系统。
主要设计目标如下:
- 以**时间复杂度为O(1)**的方式提供消息持久化能力,即使对TB级以上数据也能保证常量级复杂度的性能
- **高吞吐率,**即使在非常廉价的商用机器上也能做到单机支持每秒100K条以上消息的传输
- 支持Kafka Server间的消息分区及分布式消费,同时保证每个Partition内的消息顺序传输
- 同时支持离线数据处理和实时数据处理
- **Scale out:**支持在线水平扩展
用人体来类比,消息队列(MQ)就是数字世界的血液输送网络,而Kafka则是其中的"超级主动脉"。
表格 还在加载中,请等待加载完成后再尝试复制
Kafka的官网:http://kafka.apache.org


相关术语

为什么会有 Zookeeper Cluster 集群?
在 Kafka 架构中,ZooKeeper 是 “集群大脑”,负责元数据存储、控制器选举、Leader 选举、消费组协调等核心分布式协调工作。
三大主体
Brokers(代理)
- Kafka 集群中的服务器节点,负责存储和处理消息。
Producers(生产者)
- 向 Topics 发送消息的生产者。
- 可指定消息发送到哪个分区。(es.yml、jvm)
Consumers(消费者)
- 从 Topics 读取消息的消费者。
- 通过消费者组(Consumer Group)实现多个消费者并行处理。
三大内容
Topics(主题)
- 消息的分类单元,类似于数据库中的表。
- 每个 Topic 可分为多个分区**(Partitions)**,支持并行处理。
Partitions(分区)
- 消息保存在 Topic 中,而为了能够实现大数据的存储,一个 topic 划分为多个分区,每个分区对应一个文件,可以分别存储到不同的机器上,以实现分布式的集群存储。
- 总结起来就是,一个 topic 对应的多个 partition 分散存储到集群中的多个 broker 上,存储方式是一个 partition 对应一个文件,每个 broker 负责存储在自己机器上的 partition 中的消息读写。
- 另外,每个 partition 可以有一定的副本,备份到多台机器上,以提高可用性。
为什么要分区呢?
最根本的原因就是:Kafka基于文件进行存储,当文件内容大到一定程度时,很容易达到单个磁盘的上限。
Replica(副本)
- Kafka 还可以配置 partitions 需要备份的个数(replicas),每个 partition 将会被备份到多台机器上,以提高可用性,备份的数量可以通过配置文件指定。
- Kafka 定义了两类副本:领导者副本(Leader Replica) 和 追随者副本(Follower Replica);前者对外提供服务,后者只是被动跟随。
- 这种冗余备份的方式在分布式系统中是很常见的,那么既然有副本,就涉及到对同一个文件的多个备份如何进行管理和调度。Kafka 采取的方案是:每个 partition 选举一个 server 作为“leader”,由 leader 负责所有对该分区的读写,其他 server 作为 follower 只需要简单的与 leader 同步,保持跟进即可。如果原来的 leader 失效,会重新选举由其他的 follower 来成为新的 leader。
- 至于如何选取 leader,实际上如果我们了解 ZooKeeper,就会发现其实这正是 Zookeeper 所擅长的,Kafka 使用 ZK 在 Broker 中选出一个 Controller,用于 Partition 分配和 Leader 选举。
- 另外,这里我们可以看到,实际上作为 leader 的 server 承担了该分区所有的读写请求,因此其压力是比较大的,从整体考虑,有多少个 partition 就意味着会有多少个leader,kafka 会将 leader 分散到不同的 broker 上,确保整体的负载均衡。

先有Topic,数据量越来越大,有了Partition的概念,集群规模越大,又有了Replica副本的概念。
消息传递机制
Kafka 支持 3 种消息投递语义。在业务中,常常都是使用 At least once 的模型。
- **At most once:**最多一次,消息可能会丢失,但不会重复。
- **At least once:**最少一次,消息不会丢失,可能会重复。
- **Exactly once:**仅仅一次,消息不丢失不重复,只且消费一次。
应用场景
Kafka 的使用场景
- **活动跟踪:**Kafka 可以用来跟踪用户行为,比如我们经常回去淘宝购物,你打开淘宝的那一刻,你的登陆信息,登陆次数都会作为消息传输到 Kafka ,当你浏览购物的时候,你的浏览信息,你的搜索指数,你的购物爱好都会作为一个个消息传递给 Kafka ,这样就可以生成报告,可以做智能推荐,购买喜好等。
- **传递消息:**Kafka 另外一个基本用途是传递消息,应用程序向用户发送通知就是通过传递消息来实现的,这些应用组件可以生成消息,而不需要关心消息的格式,也不需要关心消息是如何发送的。
- **度量指标:**Kafka也经常用来记录运营监控数据。包括收集各种分布式应用的数据,生产各种操作的集中反馈,比如报警和报告。
- **日志记录:**Kafka 的基本概念来源于提交日志,比如我们可以把数据库的更新发送到 Kafka 上,用来记录数据库的更新时间,通过kafka以统一接口服务的方式开放给各种consumer,例如hadoop、Hbase、Solr等。
- **流式处理:**流式处理是有一个能够提供多种应用程序的领域。
- **限流削峰:**Kafka 多用于互联网领域某一时刻请求特别多的情况下,可以把请求写入Kafka 中,避免直接请求后端程序导致服务崩溃。
Kafka 单机部署
环境准备
准备一台虚拟机,CentOS Stream 9
| 序号 | 主机IP | 主机名 | 备注 |
|---|---|---|---|
| 1 | 192.168.88.101 | kafka1.itcast.cn | 单机模式 |
hostnamectl set-hostname kafka1.itcast.cn && bash
CentOS Stream 9 阿里源
cat >/etc/yum.repos.d/aliyun.repo<<EOF
[baseos]
name=CentOS Stream \$releasever - BaseOS
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/BaseOS/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=1
[baseos-debug]
name=CentOS Stream \$releasever - BaseOS - Debug
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/BaseOS/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[baseos-source]
name=CentOS Stream \$releasever - BaseOS - Source
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/BaseOS/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[appstream]
name=CentOS Stream \$releasever - AppStream
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/AppStream/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=1
[appstream-debug]
name=CentOS Stream \$releasever - AppStream - Debug
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/AppStream/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[appstream-source]
name=CentOS Stream \$releasever - AppStream - Source
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/AppStream/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[crb]
name=CentOS Stream \$releasever - CRB
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/CRB/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[crb-debug]
name=CentOS Stream \$releasever - CRB - Debug
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/CRB/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[crb-source]
name=CentOS Stream \$releasever - CRB - Source
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/CRB/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[highavailability]
name=CentOS Stream \$releasever - HighAvailability
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/HighAvailability/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[highavailability-debug]
name=CentOS Stream \$releasever - HighAvailability - Debug
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/HighAvailability/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[highavailability-source]
name=CentOS Stream \$releasever - HighAvailability - Source
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/HighAvailability/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[nfv]
name=CentOS Stream \$releasever - NFV
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/NFV/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[nfv-debug]
name=CentOS Stream \$releasever - NFV - Debug
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/NFV/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[nfv-source]
name=CentOS Stream \$releasever - NFV - Source
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/NFV/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[rt]
name=CentOS Stream \$releasever - RT
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/RT/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[rt-debug]
name=CentOS Stream \$releasever - RT - Debug
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/RT/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[rt-source]
name=CentOS Stream \$releasever - RT - Source
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/RT/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[resilientstorage]
name=CentOS Stream \$releasever - ResilientStorage
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/ResilientStorage/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[resilientstorage-debug]
name=CentOS Stream \$releasever - ResilientStorage - Debug
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/ResilientStorage/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[resilientstorage-source]
name=CentOS Stream \$releasever - ResilientStorage - Source
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/ResilientStorage/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[extras-common]
name=CentOS Stream \$releasever - Extras packages
baseurl=http://mirrors.aliyun.com/centos-stream/SIGs/\$stream/extras/\$basearch/extras-common/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-CentOS-SIG-Extras-SHA512
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=1
[extras-common-source]
name=CentOS Stream \$releasever - Extras packages - Source
baseurl=http://mirrors.aliyun.com/centos-stream/SIGs/\$stream/extras/source/extras-common/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-CentOS-SIG-Extras-SHA512
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
EOF
#!/bin/bash
IP_ADDRESS="192.168.88.101" # 改成自己的
NEW_HOSTNAME="kafka1.itcast.cn" # 改成自己的
hostnamectl set-hostname $NEW_HOSTNAME
source ~/.bashrc
echo "$IP_ADDRESS $NEW_HOSTNAME" >> /etc/hosts
# 关闭防火墙
iptables -F
systemctl disable firewalld
systemctl stop firewalld
# 临时关闭 SELinux
setenforce 0
# 永久关闭 SELinux
sed -i 's/SELINUX=enforcing/SELINUX=disabled/g' /etc/selinux/config
# 配置时间同步
dnf install -y chrony
systemctl enable --now chronyd
chronyc makestep
# 安装必备软件
yum install -y vim wget rsync net-tools bash-completion
# 切到新的Shell环境
bash
安装 Java
# 移除默认的 Java,防止冲突
dnf list installed | grep java
dnf remove java-1.8.0-openjdk-headless.x86_64 -y
# 安装 Java 17
dnf install java-17-openjdk-devel -y
java -version # 验证安装

为什么要安装 Java?
安装 Kafka 时需要先安装 Java,主要是因为 Kafka 是用 Scala 和 Java 语言开发的,其运行依赖于 Java 虚拟机(JVM)。
版本选择
- Kafka 2.8.0 及以上版本:建议使用 Java 8 或 11
- Kafka 3.0.0 及以上版本:建议使用 Java 11 或 17
若 Java 版本不匹配,可能导致启动失败或运行不稳定。
下载解压

版本不选最新的,一般是从最新的版本,往前找两个版本。

# 直接下载,或者用老师给的离线包进行解压
wget https://dlcdn.apache.org/kafka/3.9.0/kafka_2.13-3.9.0.tgz
tar xvf kafka_2.13-3.9.0.tgz
cd kafka_2.13-3.9.0


启动服务
# 启动 ZooKeeper(单机内置),在一个Linux终端执行
[root@kafka1 ~]# /root/kafka_2.13-3.9.0/bin/zookeeper-server-start.sh /root/kafka_2.13-3.9.0/config/zookeeper.properties
ZooKeeper 是一个分布式的、开源的协调服务,它像是一个为分布式系统设计的“管理员”或“总控中心”。
# 启动 Kafka Broker,新开一个Linux终端执行(注意是 kafka-server-start.sh,不是 zookeeper-server-start.sh)
[root@kafka1 ~]# /root/kafka_2.13-3.9.0/bin/kafka-server-start.sh /root/kafka_2.13-3.9.0/config/server.properties
Kafka Broker 就是 Kafka 消息系统的服务器节点。
你可以把它理解为:
邮局(类比于邮件系统)
货仓(类比于物流系统)
数据库服务器(类比于数据库系统)
只不过它存储和传输的是消息流(Message Streams)。
# 确认进程 注意,jps需要安装 Java 17 版本
[root@kafka1 ~]# jps | grep -E 'Kafka|QuorumPeerMain'
# 查看服务端口
ZooKeeper
[root@kafka1 ~]# netstat -pantul|grep 2181
Kafka Broker
[root@kafka1 ~]# netstat -pantul|grep 9092




[root@kafka1 ~]# jps
48422 QuorumPeerMain # ← ZooKeeper 进程
51564 Kafka # ← Kafka Broker 进程
52204 Jps # ← jps 命令自身的临时进程

jps是什么?
jps 是 JDK 提供的用于查看当前系统中所有 JVM 进程的工具,属于开发工具类。
仅包含在 openjdk-devel 包中(JRE 不包含)。
CentOS Stream 9 中通过 dnf 安装的 java-xx-openjdk-devel 会自动配置环境变量,一般无需手动操作。
可能会遇到的问题:
ZooKeeper 启动失败

错误分析
错误信息显示:
ERROR Exiting JVM with code 1
这表示 ZooKeeper 启动失败,Java 虚拟机以错误代码 1 退出。
可能的原因和解决方案
1. 端口被占用
最常见的可能是 2181 端口已经被占用:
# 检查2181端口是否被占用netstat -tlnp | grep :2181
# 或者使用 lsoflsof -i :2181
2. 数据目录问题
ZooKeeper 需要写入数据目录,可能权限不足或目录不存在:
# 检查zookeeper.properties中的dataDir配置
grep dataDir /root/kafka_2.13-3.9.0/config/zookeeper.properties
# 通常默认是:dataDir=/tmp/zookeeper# 检查该目录是否存在,是否有写入权限ls -ld /tmp/zookeeper
3. 查看详细日志
我们需要看完整的错误日志:
# 查看ZooKeeper启动日志cat /root/kafka_2.13-3.9.0/logs/zookeeper.out
# 或者查看最新的日志文件ls -la /root/kafka_2.13-3.9.0/logs/
逐步解决方案
步骤1:停止所有相关进程
# 停止可能存在的ZooKeeper进程
/root/kafka_2.13-3.9.0/bin/zookeeper-server-stop.sh
# 停止可能存在的Kafka进程
/root/kafka_2.13-3.9.0/bin/kafka-server-stop.sh
# 确认没有相关进程
jps
步骤2:检查并清理端口
# 检查2181端口占用情况
netstat -tlnp | grep :2181
# 如果有进程占用,记录PID并杀掉
# 假设PID是12345,则执行 kill -9 12345
# 对于生产环境,也可以考虑修改端口号
# 编辑zookeeper.properties,修改:# clientPort=2182
步骤3:清理和准备数据目录
# 查看配置的数据目录
grep dataDir /root/kafka_2.13-3.9.0/config/zookeeper.properties
# 假设是/tmp/zookeeper,清理并重建
rm -rf /tmp/zookeeper
mkdir -p /tmp/zookeeper
chmod 755 /tmp/zookeeper
# 如果配置的是其他目录,按相同方式处理
步骤4:以非守护进程模式启动(查看详细错误)
# 先不要用-daemon,直接在终端启动,这样可以看到完整错误
/root/kafka_2.13-3.9.0/bin/zookeeper-server-start.sh /root/kafka_2.13-3.9.0/config/zookeeper.properties
这会实时输出日志,你应该能看到具体的错误原因。
通过配置文件查看服务端口
Kafka 9092
[root@kafka1 ~]# grep 9092 /root/kafka_2.13-3.9.0/config/server.properties
# java.net.InetAddress.getCanonicalHostName(), with PLAINTEXT listener name, and
# listeners = PLAINTEXT://your.host.name:9092
#listeners=PLAINTEXT://:9092
#advertised.listeners=PLAINTEXT://your.host.name:9092
[root@kafka1 ~]#
根据配置文件,可以看出,默认端口是9092

# Kafka 默认端口是9092
[root@kafka1 ~]# ss -lnutp

Zookeeper 2181
[root@kafka1 ~]# grep 2181 /root/kafka_2.13-3.9.0/config/server.properties
zookeeper.connect=localhost:2181


Kafka 快速入门验证测试
1. 创建你的第一个主题(Topic)
主题(Topic)是什么?
就像邮箱系统中的收件箱,用于分类存储消息。
# 创建一个名为 test-topic 的主题
/root/kafka_2.13-3.9.0/bin/kafka-topics.sh \
--create \
--topic test-topic \
--bootstrap-server localhost:9092
命令解释:
--create:创建一个新主题--topic test-topic:主题名称--bootstrap-server localhost:9092:连接到本地的 Kafka 服务
成功提示:Created topic test-topic.


2. 查看所有主题
# 查看当前有哪些主题
/root/kafka_2.13-3.9.0/bin/kafka-topics.sh \
--list \
--bootstrap-server localhost:9092
预期输出:
test-topic

3. 查看主题详情
# 查看 test-topic 的详细信息
/root/kafka_2.13-3.9.0/bin/kafka-topics.sh \
--describe \
--topic test-topic \
--bootstrap-server localhost:9092
预期输出:
Topic: test-topic TopicId: xxxx PartitionCount: 1 ReplicationFactor: 1 Configs:
Topic: test-topic Partition: 0 Leader: 0 Replicas: 0 Isr: 0
字段解释:
PartitionCount: 1:有 1 个分区ReplicationFactor: 1:每个分区有 1 个副本Leader: 0:分区的主副本在 Broker ID 0 上

4. 发送消息到主题(生产者)
# 启动生产者,发送消息到 test-topic
/root/kafka_2.13-3.9.0/bin/kafka-console-producer.sh \
--topic test-topic \
--bootstrap-server localhost:9092
操作步骤:
- 执行命令后,光标会闪烁等待输入
- 输入你的第一条消息:
Hello Kafka! This is my first message.
- 按回车发送
- 继续输入更多消息:
This is the second message.
Test message number three.

- 按 Ctrl + C 退出生产者
5. 接收消息(消费者)
方式一:接收所有历史消息
# 从最开始接收所有消息
/root/kafka_2.13-3.9.0/bin/kafka-console-consumer.sh \
--topic test-topic \
--bootstrap-server localhost:9092 \
--from-beginning

方式二:实时接收新消息
# 只接收启动后发送的新消息
/root/kafka_2.13-3.9.0/bin/kafka-console-consumer.sh \
--topic test-topic \
--bootstrap-server localhost:9092
操作步骤:
- 打开一个新的终端窗口
- 执行上面的命令
- 在生产者终端发送新消息,这里会实时显示



高级主题操作
创建有多个分区的主题
分区(Partition)是什么?
就像把一个大邮箱分成多个小格子,可以并行处理消息,提高效率。
# 创建有 3 个分区的主题
/root/kafka_2.13-3.9.0/bin/kafka-topics.sh \
--create \
--topic multi-partition-topic \
--bootstrap-server localhost:9092 \
--partitions 3


创建有副本的主题
副本(Replication)是什么?
为数据创建备份,防止数据丢失,提高可靠性。
# 创建有 1 个副本的主题(单节点集群只能设1)
/root/kafka_2.13-3.9.0/bin/kafka-topics.sh \
--create \
--topic safe-topic \
--bootstrap-server localhost:9092 \
--partitions 2 \
--replication-factor 1
注意:在单节点集群中,--replication-factor 最大只能为 1,否则会报错。


命令速查表
| 功能 | 命令格式 | 示例 |
|---|---|---|
| 创建主题 | kafka-topics.sh --create --topic <名称> --bootstrap-server <地址> | --topic my-topic --bootstrap-server localhost:9092 |
| 查看主题列表 | kafka-topics.sh --list --bootstrap-server <地址> | --bootstrap-server localhost:9092 |
| 查看主题详情 | kafka-topics.sh --describe --topic <名称> --bootstrap-server <地址> | --topic test-topic --bootstrap-server localhost:9092 |
| 发送消息 | kafka-console-producer.sh --topic <名称> --bootstrap-server <地址> | --topic test-topic --bootstrap-server localhost:9092 |
| 接收消息(历史) | kafka-console-consumer.sh --topic <名称> --bootstrap-server <地址> --from-beginning | --topic test-topic --bootstrap-server localhost:9092 --from-beginning |
| 接收消息(实时) | kafka-console-consumer.sh --topic <名称> --bootstrap-server <地址> | --topic test-topic --bootstrap-server localhost:9092 |
常见选项详解
kafka-topics.sh 常用选项:
--partitions N:设置分区数量(默认为1)--replication-factor N:设置副本数量(不能超过 Broker 数量)--config <设置>:自定义主题配置
kafka-console-producer.sh 常用选项:
- 直接输入消息,一行一条
- 按
Ctrl + C退出
kafka-console-consumer.sh 常用选项:
--from-beginning:从第一条消息开始消费--partition N:只消费指定分区的消息--group <组名>:指定消费者组(实现负载均衡)
完整测试流程
# 1. 创建测试主题
/root/kafka_2.13-3.9.0/bin/kafka-topics.sh \
--create --topic demo-topic \
--bootstrap-server localhost:9092 \
--partitions 3
# 2. 查看主题详情
/root/kafka_2.13-3.9.0/bin/kafka-topics.sh \
--describe --topic demo-topic \
--bootstrap-server localhost:9092
# 3. 在终端1启动生产者
/root/kafka_2.13-3.9.0/bin/kafka-console-producer.sh \
--topic demo-topic \
--bootstrap-server localhost:9092
# 输入:Message 1、Message 2、Message 3
# 4. 在终端2启动消费者(接收历史消息)
/root/kafka_2.13-3.9.0/bin/kafka-console-consumer.sh \
--topic demo-topic \
--bootstrap-server localhost:9092 \
--from-beginning
# 5. 在终端3启动消费者(只接收新消息)
/root/kafka_2.13-3.9.0/bin/kafka-console-consumer.sh \
--topic demo-topic \
--bootstrap-server localhost:9092




实用小贴士
- 两个服务都要运行:确保 ZooKeeper 和 Kafka 都在运行
- 端口要正确:Kafka 用 9092,ZooKeeper 用 2181
- 创建主题前先检查:使用
--list查看是否已存在 - 分区数量:根据业务需求设置,分区越多并行能力越强
- 消费者组:同一个消费者组的多个消费者会分担消息处理
常见问题
Q1:创建主题时报错
Error while executing topic command : Replication factor: 2 larger than available brokers: 1
原因:副本数量超过了 Broker 数量
解决:设置 --replication-factor 1
Q2:消费者看不到消息
可能原因:
- 生产者还没发送消息
- 消费者命令没加
--from-beginning(看不到历史消息) - 主题名称拼写错误
Q3:命令执行失败
检查:
# 1. 确认 Kafka 服务运行
jps | grep Kafka
# 2. 确认端口监听
netstat -tlnp | grep :9092
# 3. 检查日志
tail -f /root/kafka_2.13-3.9.0/logs/server.log
数据存储
Kafka默认的消息存储位置是在 /tmp/kafka-logs,会被系统定期清理



可以优化相关配置
# 创建持久化目录
mkdir -p /var/lib/kafka/data
# 精确匹配 log.dirs 行进行替换
sed -i '/^log.dirs=/s|/tmp/kafka-logs|/var/lib/kafka/data|' /root/kafka_2.13-3.9.0/config/server.properties
# 验证
grep "log.dirs" /root/kafka_2.13-3.9.0/config/server.properties
/tmp/kafka-logs目录名中的 logs容易让人误以为是**“日志文件”**(如 server.log),但实际是 Kafka 的核心数据存储位置。
Kafka 将所有消息(包括 Topic、Partition 数据)以 顺序追加的日志文件(Commit Log) 形式存储在此目录,因此命名为 logs。
表格 还在加载中,请等待加载完成后再尝试复制
停止服务
先停 Kafka,再停 ZooKeeper,最后用jps验证。

# 停止 Kafka 服务
[root@kafka1 ~]# /root/kafka_2.13-3.9.0/bin/kafka-server-stop.sh
# 停止 ZooKeeper 服务
[root@kafka1 ~]# /root/kafka_2.13-3.9.0/bin/zookeeper-server-stop.sh
# 使用 jps 验证(需要满足 Java 17)
[root@kafka1 ~]# jps

启动服务
/root/kafka_2.13-3.9.0/bin/zookeeper-server-start.sh /root/kafka_2.13-3.9.0/config/zookeeper.properties
/root/kafka_2.13-3.9.0/bin/kafka-server-start.sh /root/kafka_2.13-3.9.0/config/server.properties

Kafka 集群部署
环境准备
准备三台虚拟机,CentOS Stream 9
| 序号 | 主机IP | 主机名 | 备注 |
|---|---|---|---|
| 1 | 192.168.88.101 | kafka1.itcast.cn | 节点1 |
| 2 | 192.168.88.102 | kafka2.itcast.cn | 节点2 |
| 3 | 192.168.88.103 | kafka3.itcast.cn | 节点3 |
CentOS Stream 9 阿里源
cat >/etc/yum.repos.d/aliyun.repo<<EOF
[baseos]
name=CentOS Stream \$releasever - BaseOS
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/BaseOS/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=1
[baseos-debug]
name=CentOS Stream \$releasever - BaseOS - Debug
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/BaseOS/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[baseos-source]
name=CentOS Stream \$releasever - BaseOS - Source
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/BaseOS/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[appstream]
name=CentOS Stream \$releasever - AppStream
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/AppStream/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=1
[appstream-debug]
name=CentOS Stream \$releasever - AppStream - Debug
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/AppStream/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[appstream-source]
name=CentOS Stream \$releasever - AppStream - Source
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/AppStream/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[crb]
name=CentOS Stream \$releasever - CRB
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/CRB/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[crb-debug]
name=CentOS Stream \$releasever - CRB - Debug
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/CRB/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[crb-source]
name=CentOS Stream \$releasever - CRB - Source
baseurl=https://mirrors.aliyun.com/centos-stream/\$stream/CRB/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[highavailability]
name=CentOS Stream \$releasever - HighAvailability
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/HighAvailability/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[highavailability-debug]
name=CentOS Stream \$releasever - HighAvailability - Debug
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/HighAvailability/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[highavailability-source]
name=CentOS Stream \$releasever - HighAvailability - Source
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/HighAvailability/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[nfv]
name=CentOS Stream \$releasever - NFV
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/NFV/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[nfv-debug]
name=CentOS Stream \$releasever - NFV - Debug
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/NFV/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[nfv-source]
name=CentOS Stream \$releasever - NFV - Source
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/NFV/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[rt]
name=CentOS Stream \$releasever - RT
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/RT/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[rt-debug]
name=CentOS Stream \$releasever - RT - Debug
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/RT/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[rt-source]
name=CentOS Stream \$releasever - RT - Source
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/RT/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[resilientstorage]
name=CentOS Stream \$releasever - ResilientStorage
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/ResilientStorage/\$basearch/os/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=0
[resilientstorage-debug]
name=CentOS Stream \$releasever - ResilientStorage - Debug
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/ResilientStorage/\$basearch/debug/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[resilientstorage-source]
name=CentOS Stream \$releasever - ResilientStorage - Source
baseurl=http://mirrors.aliyun.com/centos-stream/\$stream/ResilientStorage/source/tree/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-centosofficial
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
[extras-common]
name=CentOS Stream \$releasever - Extras packages
baseurl=http://mirrors.aliyun.com/centos-stream/SIGs/\$stream/extras/\$basearch/extras-common/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-CentOS-SIG-Extras-SHA512
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
countme=1
enabled=1
[extras-common-source]
name=CentOS Stream \$releasever - Extras packages - Source
baseurl=http://mirrors.aliyun.com/centos-stream/SIGs/\$stream/extras/source/extras-common/
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-CentOS-SIG-Extras-SHA512
gpgcheck=1
repo_gpgcheck=0
metadata_expire=6h
enabled=0
EOF
基础配置(三台Linux服务器都要执行)
#!/bin/bash
IP_ADDRESS="192.168.88.101" # 改成自己的
NEW_HOSTNAME="kafka1.itcast.cn" # 改成自己的
hostnamectl set-hostname $NEW_HOSTNAME
source ~/.bashrc
echo "192.168.88.101 kafka1.itcast.cn kafka1" >> /etc/hosts # 改成自己的
echo "192.168.88.102 kafka2.itcast.cn kafka2" >> /etc/hosts # 改成自己的
echo "192.168.88.103 kafka3.itcast.cn kafka3" >> /etc/hosts # 改成自己的
# 关闭防火墙
iptables -F
systemctl disable firewalld
systemctl stop firewalld
# 临时关闭 SELinux
setenforce 0
# 永久关闭 SELinux
sed -i 's/SELINUX=enforcing/SELINUX=disabled/g' /etc/selinux/config
# 配置时间同步
dnf install -y chrony
systemctl enable --now chronyd
chronyc makestep
# 安装必备软件
yum install -y vim wget rsync net-tools bash-completion
# 切换到新的Shell环境
bash
配置kafka1到各节点SSH公钥免密码连接(可选)
[root@kafka1 ~]# ssh-keygen
Generating public/private rsa key pair.
Enter file in which to save the key (/root/.ssh/id_rsa):
Enter passphrase (empty for no passphrase):
Enter same passphrase again:
Your identification has been saved in /root/.ssh/id_rsa
Your public key has been saved in /root/.ssh/id_rsa.pub
The key fingerprint is:
SHA256:A2gX5FBIHpn7WvXo1k46TAg9NCtXfJ6APrqnENkPEas root@kafka1.itcast.cn
The key's randomart image is:
+---[RSA 3072]----+
| .=*+o |
| .+B+.+ . |
| *=++ + . |
| =+oB.. o |
| E o* =So |
| ..o+ o.. |
| . +.+ .. |
| .o . =o. |
| .o ..o. |
+----[SHA256]-----+
[root@kafka1 ~]# ssh-copy-id kafka1
/usr/bin/ssh-copy-id: INFO: Source of key(s) to be installed: "/root/.ssh/id_rsa.pub"
The authenticity of host 'kafka1 (192.168.88.101)' can't be established.
ED25519 key fingerprint is SHA256:nW25vf1GUfW59rm+9bdZn6JQIgid0wPL5UEO7TR7crQ.
This host key is known by the following other names/addresses:
~/.ssh/known_hosts:1: 192.168.88.102
~/.ssh/known_hosts:4: 192.168.88.103
~/.ssh/known_hosts:5: kafka2
~/.ssh/known_hosts:6: kafka3
Are you sure you want to continue connecting (yes/no/[fingerprint])? yes
/usr/bin/ssh-copy-id: INFO: attempting to log in with the new key(s), to filter out any that are already installed
/usr/bin/ssh-copy-id: INFO: 1 key(s) remain to be installed -- if you are prompted now it is to install the new keys
root@kafka1's password:
Number of key(s) added: 1
Now try logging into the machine, with: "ssh 'kafka1'"
and check to make sure that only the key(s) you wanted were added.
[root@kafka1 ~]# ssh-copy-id kafka2
/usr/bin/ssh-copy-id: INFO: Source of key(s) to be installed: "/root/.ssh/id_rsa.pub"
/usr/bin/ssh-copy-id: INFO: attempting to log in with the new key(s), to filter out any that are already installed
/usr/bin/ssh-copy-id: INFO: 1 key(s) remain to be installed -- if you are prompted now it is to install the new keys
root@kafka2's password:
Number of key(s) added: 1
Now try logging into the machine, with: "ssh 'kafka2'"
and check to make sure that only the key(s) you wanted were added.
[root@kafka1 ~]# ssh-copy-id kafka3
/usr/bin/ssh-copy-id: INFO: Source of key(s) to be installed: "/root/.ssh/id_rsa.pub"
/usr/bin/ssh-copy-id: INFO: attempting to log in with the new key(s), to filter out any that are already installed
/usr/bin/ssh-copy-id: INFO: 1 key(s) remain to be installed -- if you are prompted now it is to install the new keys
root@kafka3's password:
Number of key(s) added: 1
Now try logging into the machine, with: "ssh 'kafka3'"
and check to make sure that only the key(s) you wanted were added.
[root@kafka1 ~]# ssh kafka1
Activate the web console with: systemctl enable --now cockpit.socket
Last login: Fri Dec 5 15:16:36 2025 from 192.168.88.1
[root@kafka1 ~]# exit
logout
Connection to kafka1 closed.
[root@kafka1 ~]# ssh kafka2
Activate the web console with: systemctl enable --now cockpit.socket
Last login: Fri Dec 5 13:29:19 2025 from 192.168.88.1
[root@kafka2 ~]# exit
logout
Connection to kafka2 closed.
[root@kafka1 ~]# ssh kafka3
Activate the web console with: systemctl enable --now cockpit.socket
Last failed login: Fri Dec 5 15:51:41 CST 2025 from 192.168.88.101 on ssh:notty
There was 1 failed login attempt since the last successful login.
Last login: Fri Dec 5 13:29:23 2025 from 192.168.88.1
[root@kafka3 ~]# exit
logout
Connection to kafka3 closed.
[root@kafka1 ~]# ssh 192.168.88.101
The authenticity of host '192.168.88.101 (192.168.88.101)' can't be established.
ED25519 key fingerprint is SHA256:nW25vf1GUfW59rm+9bdZn6JQIgid0wPL5UEO7TR7crQ.
This host key is known by the following other names/addresses:
~/.ssh/known_hosts:1: kafka1
~/.ssh/known_hosts:4: kafka2
~/.ssh/known_hosts:5: kafka3
Are you sure you want to continue connecting (yes/no/[fingerprint])? yes
Warning: Permanently added '192.168.88.101' (ED25519) to the list of known hosts.
Activate the web console with: systemctl enable --now cockpit.socket
Last login: Wed May 20 14:58:43 2026 from 192.168.88.101
[root@kafka1 ~]# exit
logout
Connection to 192.168.88.101 closed.
[root@kafka1 ~]# ssh 192.168.88.102
The authenticity of host '192.168.88.102 (192.168.88.102)' can't be established.
ED25519 key fingerprint is SHA256:nW25vf1GUfW59rm+9bdZn6JQIgid0wPL5UEO7TR7crQ.
This host key is known by the following other names/addresses:
~/.ssh/known_hosts:1: kafka1
~/.ssh/known_hosts:4: kafka2
~/.ssh/known_hosts:5: kafka3
~/.ssh/known_hosts:6: 192.168.88.101
Are you sure you want to continue connecting (yes/no/[fingerprint])? yes
Warning: Permanently added '192.168.88.102' (ED25519) to the list of known hosts.
Activate the web console with: systemctl enable --now cockpit.socket
Last login: Wed May 20 14:58:46 2026 from 192.168.88.101
[root@kafka2 ~]# exit
logout
Connection to 192.168.88.102 closed.
[root@kafka1 ~]# ssh 192.168.88.103
The authenticity of host '192.168.88.103 (192.168.88.103)' can't be established.
ED25519 key fingerprint is SHA256:nW25vf1GUfW59rm+9bdZn6JQIgid0wPL5UEO7TR7crQ.
This host key is known by the following other names/addresses:
~/.ssh/known_hosts:1: kafka1
~/.ssh/known_hosts:4: kafka2
~/.ssh/known_hosts:5: kafka3
~/.ssh/known_hosts:6: 192.168.88.101
~/.ssh/known_hosts:7: 192.168.88.102
Are you sure you want to continue connecting (yes/no/[fingerprint])? yes
Warning: Permanently added '192.168.88.103' (ED25519) to the list of known hosts.
Activate the web console with: systemctl enable --now cockpit.socket
Last login: Wed May 20 14:58:48 2026 from 192.168.88.101
[root@kafka3 ~]# exit
logout
Connection to 192.168.88.103 closed.
[root@kafka1 ~]# ssh 192.168.88.101
Activate the web console with: systemctl enable --now cockpit.socket
Last login: Wed May 20 14:58:57 2026 from 192.168.88.101
[root@kafka1 ~]# exit
logout
Connection to 192.168.88.101 closed.
[root@kafka1 ~]# ssh 192.168.88.102
Activate the web console with: systemctl enable --now cockpit.socket
Last login: Wed May 20 14:59:02 2026 from 192.168.88.101
[root@kafka2 ~]# exit
logout
Connection to 192.168.88.102 closed.
[root@kafka1 ~]# ssh 192.168.88.103
Activate the web console with: systemctl enable --now cockpit.socket
Last login: Wed May 20 14:59:05 2026 from 192.168.88.101
[root@kafka3 ~]# exit
logout
Connection to 192.168.88.103 closed.
[root@kafka1 ~]#
安装 Java(三台Linux服务器都要安装)
# 移除默认的 Java,防止冲突
dnf list installed | grep java
dnf remove java-1.8.0-openjdk-headless.x86_64 -y
# 安装 Java 17
dnf install java-17-openjdk-devel -y
# 验证安装
java -version

安装 Kafka
第一步,在其中一台主机上,执行
# 直接下载,或者用老师给的离线包进行解压
wget https://dlcdn.apache.org/kafka/3.9.0/kafka_2.13-3.9.0.tgz
tar xvf kafka_2.13-3.9.0.tgz
mv kafka_2.13-3.9.0 /opt/kafka
第二步,复制kafka
scp -r /opt/kafka/ 192.168.88.102:/opt/
scp -r /opt/kafka/ 192.168.88.103:/opt/

配置 ZooKeeper
为什么要使用 ZooKeeper?
ZooKeeper 是 Kafka 的调度中心,负责管理集群元数据、协调 Broker 和消费者,保障高可用与一致性。
专注核心功能:Kafka 主攻高性能消息传输,将协调工作外包给 ZooKeeper。
避免重复造轮子:ZooKeeper 是成熟的分布式协调服务(基于 ZAB 协议),直接复用更可靠。
ZAB 协议 是 ZooKeeper Atomic Broadcast 的缩写,即 ZooKeeper 原子广播协议。它是 ZooKeeper 专门为实现高可用、强一致性的分布式协调服务而设计的崩溃可恢复的原子广播协议。
你可以把它理解为 ZooKeeper 的“引擎”,负责在集群各节点间可靠地同步状态和数据。
从 Kafka 3.0 开始,官方推出 KRaft 模式(Kafka Raft Metadata Mode),用内置的 Raft 分布式协议替代 ZooKeeper。
目前还在持续优化阶段。

修改配置文件
注意:所有Linux服务器都要修改zookeeper.properties配置文件
# 优化zookeeper.properties配置文件
cp /opt/kafka/config/zookeeper.properties /opt/kafka/config/zookeeper.properties.bak
grep -Ev '#|^$' /opt/kafka/config/zookeeper.properties.bak > /opt/kafka/config/zookeeper.properties
cat /opt/kafka/config/zookeeper.properties

无注释版配置(可直接复制使用)
cat > /opt/kafka/config/zookeeper.properties << 'EOF'
tickTime=2000
dataDir=/var/lib/zookeeper
clientPort=2181
maxClientCnxns=60
admin.enableServer=false
initLimit=5
syncLimit=2
server.1=kafka1:2888:3888
server.2=kafka2:2888:3888
server.3=kafka3:2888:3888
autopurge.snapRetainCount=3
autopurge.purgeInterval=24
forceSync=yes
skipACL=yes
quorumListenOnAllIPs=true
preAllocSize=65536
snapCount=100000
minSessionTimeout=4000
maxSessionTimeout=40000
globalOutstandingLimit=1000
leaderServes=no
commitLogCount=500
commitLogTime=1000
EOF
带详细中文注释的配置版
cat > /opt/kafka/config/zookeeper.properties << 'EOF'
# 基本配置
tickTime=2000 # 基础时间单元(毫秒),用于计算超时时间
dataDir=/var/lib/zookeeper # ZooKeeper数据存储目录,需提前创建并设置权限
clientPort=2181 # 客户端连接端口,Kafka Broker通过此端口连接
maxClientCnxns=60 # 单个客户端IP最大连接数,防止连接数过多
admin.enableServer=false # 是否启用AdminServer,通常生产环境禁用
# 集群配置
initLimit=5 # 集群初始化时Leader等待Follower同步的最大时间(单位:tickTime)
syncLimit=2 # Leader与Follower间心跳检测的最大延迟时间(单位:tickTime)
server.1=kafka1:2888:3888 # 集群节点1:2888用于数据同步,3888用于Leader选举
server.2=kafka2:2888:3888 # 集群节点2:同上
server.3=kafka3:2888:3888 # 集群节点3:同上,奇数节点数利于选举决策
# 性能优化
autopurge.snapRetainCount=3 # 保留的旧快照文件数量,用于数据恢复和清理
autopurge.purgeInterval=24 # 自动清理任务执行间隔(小时),清理事务日志和快照
forceSync=yes # 强制数据写入磁盘,确保数据持久性,性能略有影响
skipACL=yes # 跳过ACL检查,提高性能(如果不需要细粒度权限控制)
quorumListenOnAllIPs=true # 在所有网络接口上监听,提高网络兼容性
# 高级调优
preAllocSize=65536 # 事务日志文件预分配大小(字节),减少文件碎片
snapCount=100000 # 每提交多少次事务后创建快照,影响恢复速度
minSessionTimeout=4000 # 客户端会话最小超时时间(毫秒),默认2倍tickTime
maxSessionTimeout=40000 # 客户端会话最大超时时间(毫秒),默认20倍tickTime
# 内存与垃圾回收优化
globalOutstandingLimit=1000 # 全局待处理请求上限,防止内存溢出
leaderServes=no # Leader节点是否处理客户端请求,no可减轻Leader负载;这里有一个坑,后续会详细说明,如果不想踩坑,那么这个参数可以直接删掉!
commitLogCount=500 # 内存中保留的提交日志数量,影响内存使用
commitLogTime=1000 # 提交日志刷新到磁盘的时间间隔(毫秒)
EOF
注意:leaderServes确实是 ZooKeeper 的一个参数(属于 Apache ZooKeeper 的配置项),不过它不是常用参数,而且在实际生产中很少建议手动设置,甚至很多版本中已经不推荐使用或默认行为已优化。
表格 还在加载中,请等待加载完成后再尝试复制
server.1=kafka1:2888:3888 为什么是两个端口号?
ZooKeeper 集群通过这两个端口实现以下核心功能:
Leader 选举:
当集群启动或 Leader 故障时,各节点通过 3888 端口交换投票信息,选出新的 Leader。
数据同步:
选举完成后,Follower 节点通过 2888 端口从 Leader 同步数据,确保数据一致性。



创建 myid 文件
myid 文件是 ZooKeeper 集群的基石,它:
- 标识身份:告诉节点"我是谁"
- 建立连接:让节点知道"如何找到同伴"
- 参与选举:决定在集群中的投票权重
就像团队中的工牌:每个人有唯一的工号(myid),通过工号知道自己的岗位,也能找到其他同事的位置。没有正确的 myid,节点就像没有工牌的新人,无法融入团队工作。
每个节点不同
# 节点1
mkdir -p /var/lib/zookeeper
echo "1" | tee /var/lib/zookeeper/myid
cat /var/lib/zookeeper/myid
# 节点2
mkdir -p /var/lib/zookeeper
echo "2" | tee /var/lib/zookeeper/myid
cat /var/lib/zookeeper/myid
# 节点3
mkdir -p /var/lib/zookeeper
echo "3" | tee /var/lib/zookeeper/myid
cat /var/lib/zookeeper/myid
在 Kafka 集群中,myid 是用于标识集群中每个zookeeper的唯一 ID,其值有明确的要求和规范,具体如下:
1、myid 的值必须是非负整数(0 及以上的整数)。
- 例如 1、5、100 等,不允许使用负数(如 -1 会报错)。
- 不允许使用非整数(如字符串、小数等,会导致节点启动失败)。
2、集群内所有节点的 myid 必须唯一,不能重复。
- 若两个节点使用相同的 myid,会导致元数据一致性冲突(如节点角色识别混乱、数据同步异常),最终触发集群异常或节点无法启动。
- 例如:集群中有 3 个节点,myid 可设置为 1、2、3,但不能有两个节点同时为 1。
配置 Kafka 集群
编辑每台节点的 server.properties
# 修改 /opt/kafka/config/server.properties
vim /opt/kafka/config/server.properties
... (省略)
23 # The id of the broker. This must be set to a unique integer for each broker.
24 broker.id=1 # 节点1设为1,节点2设为2,节点3设为3
30 # FORMAT:
31 # listeners = listener_name://host_name:port
32 # EXAMPLE:
33 # listeners = PLAINTEXT://your.host.name:9092
34 listeners=PLAINTEXT://kafka1:9092 # 修改为对应节点的主机名
61 # A comma separated list of directories under whichto store log files
62 log.dirs=/opt/kafka/kafka-logs
64 # The default number of log partitions per topic. More partitions allow greater
65 # parallelism for consumption, but this will also result in more files across
66 # the brokers.
67 num.partitions=3
122 # server. e.g. "127.0.0.1:3000,127.0.0.1:3001,127.0.0.1:3002".
123 # You can also append an optional chroot string to the urls to specify the
124 # root directory for all kafka znodes.
125 zookeeper.connect=kafka1:2181,kafka2:2181,kafka3:2181
# 表示 Kafka 中自动创建主题的默认副本因子为 3。即对于自动创建的主题,其每个分区都会有 3 个副本
default.replication.factor=3
# Kafka 通过同步副本集合(In - Sync Replicas,ISR)来管理副本同步状态,只有在 ISR 中的副本才被认为是与 Leader 副本保持同步的
min.insync.replicas=2
... (省略)
注意:
上面内容不能照抄,kafka1、kafka2、kakfa3换成实际的域名地址
优化各节点的/opt/kafka/config/server.properties配置文件
cp /opt/kafka/config/server.properties /opt/kafka/config/server.properties.bak
sed -i '/^\s*#/d; /^\s*$/d' /opt/kafka/config/server.properties
cat /opt/kafka/config/server.properties
在各节点创建/opt/kafka/kafka-logs目录
mkdir -p /opt/kafka/kafka-logs




三个节点的无注释版配置文件
Kafka1 节点配置
cat > /opt/kafka/config/server.properties << 'EOF'
broker.id=1
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
listeners=PLAINTEXT://kafka1:9092
advertised.listeners=PLAINTEXT://kafka1:9092
log.dirs=/opt/kafka/kafka-logs
num.partitions=3
num.recovery.threads.per.data.dir=1
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
default.replication.factor=3
min.insync.replicas=2
log.retention.hours=168
log.retention.check.interval.ms=300000
zookeeper.connect=kafka1:2181,kafka2:2181,kafka3:2181
zookeeper.connection.timeout.ms=18000
group.initial.rebalance.delay.ms=0
EOF
Kafka2 节点配置
cat > /opt/kafka/config/server.properties << 'EOF'
broker.id=2
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
listeners=PLAINTEXT://kafka2:9092
advertised.listeners=PLAINTEXT://kafka2:9092
log.dirs=/opt/kafka/kafka-logs
num.partitions=3
num.recovery.threads.per.data.dir=1
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
default.replication.factor=3
min.insync.replicas=2
log.retention.hours=168
log.retention.check.interval.ms=300000
zookeeper.connect=kafka1:2181,kafka2:2181,kafka3:2181
zookeeper.connection.timeout.ms=18000
group.initial.rebalance.delay.ms=0
EOF
Kafka3 节点配置
cat > /opt/kafka/config/server.properties << 'EOF'
broker.id=3
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
listeners=PLAINTEXT://kafka3:9092
advertised.listeners=PLAINTEXT://kafka3:9092
log.dirs=/opt/kafka/kafka-logs
num.partitions=3
num.recovery.threads.per.data.dir=1
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
default.replication.factor=3
min.insync.replicas=2
log.retention.hours=168
log.retention.check.interval.ms=300000
zookeeper.connect=kafka1:2181,kafka2:2181,kafka3:2181
zookeeper.connection.timeout.ms=18000
group.initial.rebalance.delay.ms=0
EOF
带详细中文注释的配置文件(适用于所有节点,只需修改 broker.id 和 listeners)
Kafka 通用配置模板(带中文注释)
cat > /opt/kafka/config/server.properties << 'EOF'
# 基础配置
broker.id=1 # 代理ID,集群中必须唯一(kafka1=1, kafka2=2, kafka3=3)
num.network.threads=3 # 处理网络请求的线程数
num.io.threads=8 # 处理磁盘I/O的线程数,通常设置为磁盘数量×2
# 网络缓冲区配置
socket.send.buffer.bytes=102400 # 发送缓冲区大小(字节)
socket.receive.buffer.bytes=102400 # 接收缓冲区大小(字节)
socket.request.max.bytes=104857600 # 单个请求最大大小(100MB)
# 网络监听配置
listeners=PLAINTEXT://kafka1:9092 # 内部监听地址(kafka1/kafka2/kafka3对应修改)
advertised.listeners=PLAINTEXT://kafka1:9092 # 对外发布的地址(客户端连接使用)
# 日志存储配置
log.dirs=/opt/kafka/kafka-logs # Kafka数据日志存储目录(比/tmp更持久)
num.partitions=3 # 主题默认分区数,提高并行消费能力
# 恢复线程配置
num.recovery.threads.per.data.dir=1 # 每个数据目录的恢复线程数
# 内部主题副本配置
offsets.topic.replication.factor=3 # __consumer_offsets主题的副本因子(确保高可用)
transaction.state.log.replication.factor=3 # 事务状态日志的副本因子
transaction.state.log.min.isr=2 # 事务状态日志的最小同步副本数
# 主题默认配置
default.replication.factor=3 # 自动创建主题的默认副本数(建议3副本)
min.insync.replicas=2 # 最小同步副本数,保证数据安全
# 日志保留策略
log.retention.hours=168 # 日志保留时间(7天)
log.retention.check.interval.ms=300000 # 日志清理检查间隔(5分钟)
# ZooKeeper连接配置
zookeeper.connect=kafka1:2181,kafka2:2181,kafka3:2181 # ZooKeeper集群地址(全部节点)
zookeeper.connection.timeout.ms=18000 # ZooKeeper连接超时时间(18秒)
# 消费组配置
group.initial.rebalance.delay.ms=0 # 消费组再平衡延迟时间(0=立即开始)
EOF
各节点配置差异说明
表格 还在加载中,请等待加载完成后再尝试复制
启动集群
启动 ZooKeeper(所有节点依次执行)
# 先启动 ZooKeeper
/opt/kafka/bin/zookeeper-server-start.sh -daemon /opt/kafka/config/zookeeper.properties
# 等待3-5秒,验证测试
ss -nlutp | grep 2181 # 看一下 2181 端口是否绑定给,没绑定,别启动 Kafka,先定位 bug


测试所有 ZooKeeper 节点连接情况(以下测试命令在任意一个节点都可以执行)
for ip in 192.168.88.101 192.168.88.102 192.168.88.103; do
echo "测试 $ip:2181..."
timeout 5 /opt/kafka/bin/zookeeper-shell.sh $ip:2181 <<< "ls /" 2>&1 | grep -q "WatchedEvent state:SyncConnected" && echo "连接成功" || echo "连接失败"
done
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.101:2181 <<< "ls /"
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.102:2181 <<< "ls /"
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.103:2181 <<< "ls /"



说明:如果看到某个节点连接失败,在这里先不用管,后续会讲解具体处理方法!

启动 Kafka(所有节点依次执行)
# 启动 Kafka
/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties
# 查看日志
tail -200 /opt/kafka/logs/server.log
或者
tail -f /opt/kafka/logs/server.log
# 等待3-5秒,验证测试
ss -nlutp | grep 9092
或者
netstat -pantul | grep 9092


清理环境重启 Kafka 服务(所有节点依次执行)
/opt/kafka/bin/kafka-server-stop.sh
pkill -9 -f Kafka
rm -rf /opt/kafka/kafka-logs/*
/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties
netstat -pantul|grep 9092
验证测试 ZooKeeper
# 验证集群状态(因为是集群,在任意一个 ZooKeeper 节点上执行就可以)
/opt/kafka/bin/zookeeper-shell.sh kafka1:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka2:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka3:2181 ls /brokers/ids
结果如下:
Connecting to kafkaX:2181 # 表示客户端正在尝试连接到 ZooKeeper 服务
WATCHER:: # 表示后续内容是 ZooKeeper 的 Watcher(监听器)事件回调信息
WatchedEvent state:SyncConnected type:None path:null # 客户端与 ZooKeeper 服务器建立了同步连接
[2] # [2] 可能表示当前客户端的会话 ID 片段,或只是一个简单的序号标识

Kafka 和 ZooKeeper 故障分析与解决(难点)


使用IP地址替换主机名
首先确保所有节点的/etc/hosts解释正确!
cat >/etc/hosts<<EOF
192.168.88.101 kafka1.itcast.cn kafka1
192.168.88.102 kafka2.itcast.cn kafka2
192.168.88.103 kafka3.itcast.cn kafka3
EOF
步骤1:在所有节点上更新 zookeeper 配置
kafka1 节点配置:
cat > /opt/kafka/config/zookeeper.properties << 'EOF'
tickTime=2000
dataDir=/var/lib/zookeeper
clientPort=2181
maxClientCnxns=60
admin.enableServer=false
initLimit=5
syncLimit=2
server.1=192.168.88.101:2888:3888
server.2=192.168.88.102:2888:3888
server.3=192.168.88.103:2888:3888
autopurge.snapRetainCount=3
autopurge.purgeInterval=24
forceSync=yes
skipACL=yes
quorumListenOnAllIPs=true
preAllocSize=65536
snapCount=100000
minSessionTimeout=4000
maxSessionTimeout=40000
globalOutstandingLimit=1000
leaderServes=yes
commitLogCount=500
commitLogTime=1000
EOF
kafka2 节点配置:
cat > /opt/kafka/config/zookeeper.properties << 'EOF'
tickTime=2000
dataDir=/var/lib/zookeeper
clientPort=2181
maxClientCnxns=60
admin.enableServer=false
initLimit=5
syncLimit=2
server.1=192.168.88.101:2888:3888
server.2=192.168.88.102:2888:3888
server.3=192.168.88.103:2888:3888
autopurge.snapRetainCount=3
autopurge.purgeInterval=24
forceSync=yes
skipACL=yes
quorumListenOnAllIPs=true
preAllocSize=65536
snapCount=100000
minSessionTimeout=4000
maxSessionTimeout=40000
globalOutstandingLimit=1000
leaderServes=yes
commitLogCount=500
commitLogTime=1000
EOF
kafka3 节点配置:
cat > /opt/kafka/config/zookeeper.properties << 'EOF'
tickTime=2000
dataDir=/var/lib/zookeeper
clientPort=2181
maxClientCnxns=60
admin.enableServer=false
initLimit=5
syncLimit=2
server.1=192.168.88.101:2888:3888
server.2=192.168.88.102:2888:3888
server.3=192.168.88.103:2888:3888
autopurge.snapRetainCount=3
autopurge.purgeInterval=24
forceSync=yes
skipACL=yes
quorumListenOnAllIPs=true
preAllocSize=65536
snapCount=100000
minSessionTimeout=4000
maxSessionTimeout=40000
globalOutstandingLimit=1000
leaderServes=yes
commitLogCount=500
commitLogTime=1000
EOF
步骤2:清理所有节点并设置 myid
在 kafka1 节点执行:
# 1. 停止服务
/opt/kafka/bin/zookeeper-server-stop.sh
/opt/kafka/bin/kafka-server-stop.sh
sleep 3
pkill -9 -f "QuorumPeerMain|Kafka" 2>/dev/null
# 2. 清理并设置数据目录
rm -rf /var/lib/zookeeper/*
mkdir -p /var/lib/zookeeper
chmod 755 /var/lib/zookeeper && chown -R root:root /var/lib/zookeeper
echo "1" > /var/lib/zookeeper/myid
在 kafka2 节点执行:
# 1. 停止服务
/opt/kafka/bin/zookeeper-server-stop.sh
/opt/kafka/bin/kafka-server-stop.sh
sleep 3
pkill -9 -f "QuorumPeerMain|Kafka" 2>/dev/null
# 2. 清理并设置数据目录
rm -rf /var/lib/zookeeper/*
mkdir -p /var/lib/zookeeper
chmod 755 /var/lib/zookeeper && chown -R root:root /var/lib/zookeeper
echo "2" > /var/lib/zookeeper/myid
在 kafka3 节点执行:
# 1. 停止服务
/opt/kafka/bin/zookeeper-server-stop.sh
/opt/kafka/bin/kafka-server-stop.sh
sleep 3
pkill -9 -f "QuorumPeerMain|Kafka" 2>/dev/null
# 2. 清理并设置数据目录
rm -rf /var/lib/zookeeper/*
mkdir -p /var/lib/zookeeper
chmod 755 /var/lib/zookeeper && chown -R root:root /var/lib/zookeeper
echo "3" > /var/lib/zookeeper/myid
步骤3:按顺序启动 zookeeper 集群
1. 先启动 kafka1 的 zookeeper:
# 在 kafka1 执行
/opt/kafka/bin/zookeeper-server-start.sh -daemon /opt/kafka/config/zookeeper.properties
sleep 5
echo "检查 kafka1 zookeeper 启动状态..."
jps
2. 再启动 kafka2 的 zookeeper:
# 在 kafka2 执行
/opt/kafka/bin/zookeeper-server-start.sh -daemon /opt/kafka/config/zookeeper.properties
sleep 5
echo "检查 kafka2 zookeeper 启动状态..."
jps
3. 最后启动 kafka3 的 zookeeper:
# 在 kafka3 执行
/opt/kafka/bin/zookeeper-server-start.sh -daemon /opt/kafka/config/zookeeper.properties
sleep 5
echo "检查 kafka3 zookeeper 启动状态..."
jps
步骤4:验证 zookeeper 集群
在任意节点上测试:
echo "=== 测试所有 zookeeper 节点 ==="
for ip in 192.168.88.101 192.168.88.102 192.168.88.103; do
echo "测试 $ip:2181..."
timeout 5 /opt/kafka/bin/zookeeper-shell.sh $ip:2181 <<< "ls /" 2>&1 | grep -q "WatchedEvent state:SyncConnected" && echo "连接成功" || echo "连接失败"
done
或者
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.101:2181 <<< "ls /"
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.102:2181 <<< "ls /"
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.103:2181 <<< "ls /"
步骤5:启动 Kafka 服务
在 kafka1 启动 Kafka:
rm -rf /opt/kafka/kafka-logs/* # 清理残留数据
/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties
sleep 5
jps
在 kafka2 启动 Kafka:
rm -rf /opt/kafka/kafka-logs/* # 清理残留数据
/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties
sleep 5
jps
在 kafka3 启动 Kafka:
rm -rf /opt/kafka/kafka-logs/* # 清理残留数据
/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties
sleep 5
jps
步骤6:验证 Kafka 集群状态
在任意节点上执行:
# 验证集群状态(因为是集群,在任意一个 zookeeper 节点上执行都可以)
for host in kafka1 kafka2 kafka3
do
echo "=== 测试 $host ==="
/opt/kafka/bin/zookeeper-shell.sh $host:2181 ls /brokers/ids 2>&1 | grep -q "WatchedEvent state:SyncConnected"
if [ $? -eq 0 ]; then
echo "集群状态正常"
else
echo "集群状态异常"
fi
done
或者
/opt/kafka/bin/zookeeper-shell.sh kafka1:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka2:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka3:2181 ls /brokers/ids
echo "=== 检查 Kafka Broker 注册状态 ==="
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.101:2181 <<< "ls /brokers/ids"
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.102:2181 <<< "ls /brokers/ids"
/opt/kafka/bin/zookeeper-shell.sh 192.168.88.103:2181 <<< "ls /brokers/ids"
步骤7:创建测试主题验证功能
# 创建测试主题
/opt/kafka/bin/kafka-topics.sh --create \
--topic test-cluster \
--bootstrap-server 192.168.88.101:9092 \
--partitions 3 \
--replication-factor 3
# 查看主题列表
/opt/kafka/bin/kafka-topics.sh --list \
--bootstrap-server 192.168.88.101:9092
# 查看主题详情
/opt/kafka/bin/kafka-topics.sh --describe \
--topic test-cluster \
--bootstrap-server 192.168.88.101:9092



如果遇到主题Topic已经存在,则按照以下方法解决:
[root@kafka1 ~]# /opt/kafka/bin/kafka-topics.sh --create \
--topic test-cluster \
--bootstrap-server 192.168.88.101:9092 \
--partitions 3 \
--replication-factor 3
Error while executing topic command : Topic 'test-cluster' already exists.
[2025-12-07 10:16:22,482] ERROR org.apache.kafka.common.errors.TopicExistsException: Topic 'test-cluster' already exists.
(org.apache.kafka.tools.TopicCommand)
/opt/kafka/bin/kafka-topics.sh --delete \
--topic test-cluster \
--bootstrap-server 192.168.88.101:9092
注意主题Topic名称以及IP地址

验证集群状态
for ip in 192.168.88.101 192.168.88.102 192.168.88.103; do echo "测试 $ip:2181..."
timeout 5 /opt/kafka/bin/zookeeper-shell.sh $ip:2181 <<< "ls /" 2>&1 | grep -q "WatchedEvent state:SyncConnected" && echo "连接成功" || echo "连接失败"; done
/opt/kafka/bin/zookeeper-shell.sh kafka1:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka2:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka3:2181 ls /brokers/ids

for ip in 192.168.88.101 192.168.88.102 192.168.88.103; do
echo "=== $ip ==="
echo "srvr" | nc -w 2 $ip 2181 | grep "Mode:"
done

| 部分 | 含义 | |
|---|---|---|
echo "srvr" | 输出字符串 srvr(ZooKeeper 的四字命令) | |
| ` | ` | 管道符,将前一个命令的输出传递给后一个命令 |
nc | Netcat,网络工具界的"瑞士军刀" | |
-w 2 | 超时时间为 2 秒 | |
$ip | 目标服务器 IP 地址(变量) | |
2181 | ZooKeeper 客户端端口 |
配置 systemd 服务 zookeeper.service 和 kafka.service
停止原来的 Kafka 和 ZooKeeper 服务
/opt/kafka/bin/kafka-server-stop.sh
/opt/kafka/bin/zookeeper-server-stop.sh
查看 Kafka 和 ZooKeeper 服务状态
ps aux|grep java
netstat -pantul|grep 2181
netstat -pantul|grep 9092
创建 zookeeper.service 文件
cat > /etc/systemd/system/zookeeper.service <<EOF
[Unit]
Description=Apache Zookeeper daemon (Kafka dependency)
Documentation=https://zookeeper.apache.org/
After=network.target network-online.target
Wants=network-online.target
[Service]
Type=simple
WorkingDirectory=/opt/kafka
ExecStart=/opt/kafka/bin/zookeeper-server-start.sh /opt/kafka/config/zookeeper.properties
ExecStop=/opt/kafka/bin/zookeeper-server-stop.sh
ExecReload=/bin/kill -HUP \$MAINPID
Restart=on-failure
RestartSec=30s
TimeoutStartSec=5min
TimeoutStopSec=5min
LimitNOFILE=65535
LimitNPROC=65535
KillMode=mixed
LimitCORE=infinity
[Install]
WantedBy=multi-user.target
EOF
systemctl daemon-reload
systemctl enable zookeeper --now
systemctl status zookeeper --no-pager
创建 kafka.service 文件
cat > /etc/systemd/system/kafka.service <<EOF
[Unit]
Description=Kafka server daemon
Documentation=http://kafka.apache.org/
After=network.target network-online.target zookeeper.service
Wants=network-online.target zookeeper.service
Requires=zookeeper.service
[Service]
Type=simple
WorkingDirectory=/opt/kafka
ExecStart=/opt/kafka/bin/kafka-server-start.sh /opt/kafka/config/server.properties
ExecStop=/opt/kafka/bin/kafka-server-stop.sh
ExecReload=/bin/kill -HUP \$MAINPID
Restart=on-failure
RestartSec=30s
TimeoutStartSec=5min
TimeoutStopSec=5min
LimitNOFILE=65535
LimitNPROC=65535
KillMode=mixed
LimitCORE=infinity
[Install]
WantedBy=multi-user.target
EOF
systemctl daemon-reload
systemctl enable kafka --now
systemctl status kafka --no-pager
注意:
systemd 停止 Kafka 和 ZooKeeper 服务
要先执行 systemctl stop kafka
然后再执行 systemctl stop zookeeper
Kafka 依赖 ZooKeeper,但 ZooKeeper 不依赖 Kafka
- 先停 Kafka:确保 Kafka 这个“上层服务”先优雅退出,不再向 ZooKeeper 发起请求。
- 再停 ZooKeeper:等 Kafka 完全停掉后,再停 ZooKeeper,避免 Kafka 在运行过程中突然失去 ZooKeeper 导致异常或数据不一致。
类比理解:
就像关电脑前要先保存文档,再关机:
- Kafka 是“文档”,ZooKeeper 是“操作系统”。
- 如果先关 ZooKeeper(操作系统),Kafka(文档)可能还没保存完,数据就坏了。
Kafka+ZooKeeper 集群验证核心方法
# 修改 leaderServes=yes 为 leaderServes=no
sed -i 's/leaderServes=yes/leaderServes=no/' /opt/kafka/config/zookeeper.properties
# ========== 添加注释(禁用 leaderServes=no)==========
sed -i 's/^leaderServes=no/#leaderServes=no/' /opt/kafka/config/zookeeper.properties
cat /opt/kafka/config/zookeeper.properties
# ========== 取消注释(启用 leaderServes=no)==========
sed -i 's/^#leaderServes=no/leaderServes=no/' /opt/kafka/config/zookeeper.properties
cat /opt/kafka/config/zookeeper.properties
# 重启 zookeeper
systemctl restart zookeeper
systemctl status zookeeper --no-pager
# 检查集群状态
for ip in 192.168.88.101 192.168.88.102 192.168.88.103; do echo "测试 $ip:2181..."
timeout 5 /opt/kafka/bin/zookeeper-shell.sh $ip:2181 <<< "ls /" 2>&1 | grep -q "WatchedEvent state:SyncConnected" && echo "连接成功" || echo "连接失败"; done
# 查看 brokers 状态
/opt/kafka/bin/zookeeper-shell.sh kafka1:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka2:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka3:2181 ls /brokers/ids
# 查看 leader 与 follower
for ip in 192.168.88.101 192.168.88.102 192.168.88.103; do
echo "=== $ip ==="
echo "srvr" | nc -w 2 $ip 2181 | grep "Mode:"
done
# 查看日志
journalctl -u zookeeper | grep -i "leader\|unknown\|invalid\|warn"
| 配置 | 建议 |
|---|---|
leaderServes=no | ❌ 不要使用,已废弃,会导致连接问题 |
#leaderServes=no(注释) | ✅ 保持,使用默认行为 |
| 升级 ZooKeeper | ✅ 建议升级到 3.6+,该参数已不存在 |
Kafka+ZooKeeper 集群测试(重点)
Kafka+ZooKeeper 集群正常!
下面提供完整的消息生产/消费测试步骤,可以跨节点验证集群可用性(比如在 kafka1 生产消息,在 kafka2/kafka3 消费消息)。
测试前提
所有节点的 Kafka 进程已启动,且 ZooKeeper 中能看到 brokers/ids 包含 [1,2,3](已满足)。
# 查看 brokers 状态
/opt/kafka/bin/zookeeper-shell.sh kafka1:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka2:2181 ls /brokers/ids
/opt/kafka/bin/zookeeper-shell.sh kafka3:2181 ls /brokers/ids
for ip in 192.168.88.101 192.168.88.102 192.168.88.103; do echo "测试 $ip:2181..."
timeout 5 /opt/kafka/bin/zookeeper-shell.sh $ip:2181 <<< "ls /" 2>&1 | grep -q "WatchedEvent state:SyncConnected" && echo "连接成功" || echo "连接失败"; done
for ip in 192.168.88.101 192.168.88.102 192.168.88.103; do
echo "=== $ip ==="
echo "srvr" | nc -w 2 $ip 2181 | grep "Mode:"
done

步骤1:创建测试 Topic(任意节点执行,比如 kafka1)
先创建一个跨集群的测试 Topic(3副本,3分区,适配3节点集群):
# 在kafka1节点执行:创建topic test_topic
/opt/kafka/bin/kafka-topics.sh \
--create \
--topic test_topic \
--bootstrap-server 192.168.88.101:9092,192.168.88.102:9092,192.168.88.103:9092 \
--partitions 3 \
--replication-factor 3
# 验证topic创建成功
/opt/kafka/bin/kafka-topics.sh \
--describe \
--topic test_topic \
--bootstrap-server 192.168.88.101:9092
正常输出示例:
Topic: test_topic TopicId: xxxxx PartitionCount: 3 ReplicationFactor: 3 Configs:
Topic: test_topic Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
Topic: test_topic Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
Topic: test_topic Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2

步骤2:在 kafka1 节点启动「消息生产者」(发送测试消息)
# 在kafka1节点执行:启动生产者,输入消息后按回车发送(可多行输入)
/opt/kafka/bin/kafka-console-producer.sh \
--topic test_topic \
--bootstrap-server 192.168.88.101:9092
# 输入测试消息(示例,可自定义):
Hello Kafka Cluster
This is a test message from kafka1
Cluster is working normally
输入完成后,按 Ctrl+C 退出生产者。

步骤3:在 kafka2/kafka3 节点启动「消息消费者」(跨节点消费)
方式1:在 kafka2 节点消费(实时监听)
# 登录kafka2节点执行:启动消费者,从头消费所有消息
/opt/kafka/bin/kafka-console-consumer.sh \
--topic test_topic \
--bootstrap-server 192.168.88.102:9092 \
--from-beginning
预期结果:会看到 kafka1 发送的所有测试消息:
Hello Kafka Cluster
This is a test message from kafka1
Cluster is working normally


方式2:在 kafka3 节点消费(验证多节点一致性)
# 登录kafka3节点执行:同样消费消息
/opt/kafka/bin/kafka-console-consumer.sh \
--topic test_topic \
--bootstrap-server 192.168.88.103:9092 \
--from-beginning
预期结果:和 kafka2 节点看到的消息完全一致,证明集群消息同步正常。

步骤4:进阶测试(可选)
1. 多生产者/多消费者测试
- 在 kafka2 节点也启动生产者,发送消息:
/opt/kafka/bin/kafka-console-producer.sh --topic test_topic --bootstrap-server 192.168.88.102:9092
# 输入:message from kafka2
- 在 kafka1 节点启动消费者,能看到 kafka2 发送的消息,验证双向通信。
/opt/kafka/bin/kafka-console-consumer.sh \
--topic test_topic \
--bootstrap-server 192.168.88.101:9092 \
--from-beginning

2. 验证 Topic 副本同步
# 在任意节点执行,查看topic副本状态
/opt/kafka/bin/kafka-topics.sh \
--describe \
--topic test_topic \
--bootstrap-server 192.168.88.101:9092
重点看 Isr 列(In-Sync Replicas),需包含 1,2,3,说明所有副本同步正常。

测试清理(可选)
若测试完成后需要删除测试 Topic:
/opt/kafka/bin/kafka-topics.sh \
--delete \
--topic test_topic \
--bootstrap-server 192.168.88.101:9092
核心验证点
- 跨节点生产/消费消息正常 → 集群网络和副本同步正常;
- 所有节点的消费者都能读到相同消息 → ZooKeeper 元数据和 Kafka 数据同步正常;
Isr包含所有节点 → 副本健康,集群高可用。
至此, Kafka+ZooKeeper 集群已完全具备生产可用的条件!
验证 Kafka+ZooKeeper 集群容错性(重点)
# 创建 Topic(3副本)
/opt/kafka/bin/kafka-topics.sh --create --topic cluster-test \
--bootstrap-server kafka1:9092 \
--partitions 3 \
--replication-factor 3
# 查看 topic
/opt/kafka/bin/kafka-topics.sh \
--describe --topic cluster-test \
--bootstrap-server kafka1:9092
# 模拟节点宕机(关闭节点3)
ssh kafka3 "pkill -f kafka.Kafka" # kafka.Kafka 是 Kafka 主类的名称,通常出现在启动命令中
或者
ssh kafka3 "systemctl stop kafka"
# 如果杀不掉,进到对应主机,kill
kill -9 进程的PID(jps查看)
# 查看 Topic 状态(应显示 Leader 切换)
/opt/kafka/bin/kafka-topics.sh --describe --topic cluster-test --bootstrap-server kafka1:9092
# 启动 kafka3
ssh kafka3 "/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties"
或者
ssh kafka3 "systemctl start kafka"


验证结果完全符合 Kafka 集群高可用的预期,整个故障转移过程正常。
下面逐行拆解分析,理解每个字段的含义和集群的行为逻辑:
一、初始状态(kafka3 未宕机)
Topic: cluster-test PartitionCount: 3 ReplicationFactor: 3 Configs: min.insync.replicas=2
Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2
核心字段解读
| 字段 | 含义 |
|---|---|
Partition | 分区编号(0/1/2,共3个分区,与你指定的 --partitions 3 一致) |
Leader | 该分区的主副本(处理读写请求的节点): - 分区0主节点:kafka1(broker.id=1) - 分区1主节点:kafka2(broker.id=2) - 分区2主节点:kafka3(broker.id=3) |
Replicas | 该分区的所有副本列表(3副本,与 --replication-factor 3 一致) |
Isr | 同步副本列表(In-Sync Replicas):所有副本都同步完成,无延迟 |
min.insync.replicas=2 | 最小同步副本数:至少2个副本同步,才允许写入(默认值,保障数据可靠性) |
初始状态结论
- 集群3节点(kafka1/2/3)全部正常,所有分区的主副本、副本列表、同步副本均符合预期;
- 数据可正常读写,集群处于完全健康状态。
二、kafka3 宕机后状态
Topic: cluster-test PartitionCount: 3 ReplicationFactor: 3 Configs: min.insync.replicas=2
Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2
Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,1
Partition: 2 Leader: 1 Replicas: 3,1,2 Isr: 1,2
核心变化分析
1. Isr 列表变化:所有分区的 Isr 都移除了 3(kafka3)
- 原因:kafka3 宕机后,无法再同步数据,Kafka 会将其从「同步副本列表」中剔除;
- 结果:所有分区的 Isr 剩余 2 个节点(满足
min.insync.replicas=2),集群仍可正常写入。
2. 分区2的 Leader 从 3 切换为 1(核心故障转移)
- 原状态:分区2的主节点是 kafka3(broker.id=3),宕机后主节点不可用;
- 故障转移:Kafka 控制器(Controller)自动从该分区的 Isr 列表(3,1,2)中选举新主节点 → 选举 kafka1(broker.id=1);
- 意义:分区2的读写请求无缝切换到 kafka1,业务无感知,实现高可用。
3. 其他分区 Leader 未变化(符合预期)
- 分区0主节点是 kafka1(正常),无需切换;
- 分区1主节点是 kafka2(正常),无需切换;
- 仅宕机节点作为主节点的分区(分区2)触发 Leader 切换,资源消耗最小。
三、整体结论
你的 Kafka 集群高可用机制完全生效,验证结果符合预期:
-
宕机前:集群全节点正常,副本同步完整;
-
宕机后:
-
- 自动剔除宕机节点的同步副本(Isr 列表更新);
- 自动选举新主节点(分区2 Leader 切换);
- 剩余节点满足最小同步副本数,集群仍可正常读写;
-
无异常行为:未出现分区不可用、副本丢失、同步异常等问题。
四、测试产生消息与消费消息(这部分测试可以暂时先放一放)
# 删除test_topic666主题,清理环境
/opt/kafka/bin/kafka-topics.sh \
--delete \
--topic test_topic666 \
--bootstrap-server 192.168.88.101:9092
# 在kafka1节点执行:创建topic test_topic666
/opt/kafka/bin/kafka-topics.sh \
--create \
--topic test_topic666 \
--bootstrap-server 192.168.88.101:9092,192.168.88.102:9092 \
--partitions 2 \
--replication-factor 2
或者
# 如果kafka3已经修复
/opt/kafka/bin/kafka-topics.sh \
--create \
--topic test_topic666 \
--bootstrap-server 192.168.88.101:9092,192.168.88.102:9092,192.168.88.103:9092 \
--partitions 3 \
--replication-factor 3
# 验证topic创建成功
/opt/kafka/bin/kafka-topics.sh \
--describe \
--topic test_topic \
--bootstrap-server 192.168.88.101:9092
# 在kafka1节点执行:启动生产者,输入消息后按回车发送(可多行输入)
/opt/kafka/bin/kafka-console-producer.sh \
--topic test_topic \
--bootstrap-server 192.168.88.101:9092
# 输入测试消息(示例,可自定义):
Hello Kafka Cluster
This is a test message from kafka1
Cluster is working normally
# 登录kafka2节点执行:启动消费者,从头消费所有消息
/opt/kafka/bin/kafka-console-consumer.sh \
--topic test_topic \
--bootstrap-server 192.168.88.102:9092 \
--from-beginning
补充:恢复 kafka3 后的预期行为(可选验证)
若重启 kafka3,执行 ssh kafka3 "/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties",再查看 Topic 状态:
- kafka3 会自动追上缺失的同步数据;
- 所有分区的 Isr 列表会重新加入 3;
- 分区2的 Leader 会切回 3(可选,默认 Kafka 不会主动切回,需手动触发或等待参数配置)。

手动切回分区 2 的 Leader 为 3(可选)
一、新版本(Kafka 3.0+)优先副本选举命令
# 直接执行优先副本选举(无需手动创建JSON文件,指定topic即可)
/opt/kafka/bin/kafka-leader-election.sh \
--bootstrap-server kafka1:9092 \
--topic cluster-test \
--partition 2 \
--election-type preferred
二、命令解读
| 参数 | 含义 |
|---|---|
--bootstrap-server | 指定Kafka集群地址(任意可用节点即可) |
--topic | 要操作的Topic名称(cluster-test) |
--partition | 要切换Leader的分区编号(2) |
--election-type preferred | 选举类型:优先副本(即恢复到Replicas列表中第一个节点作为Leader) |
三、执行后验证
/opt/kafka/bin/kafka-topics.sh --describe --topic cluster-test --bootstrap-server kafka1:9092
预期结果:分区2的 Leader 会切回 3(kafka3),Isr 仍为 1,2,3:
Topic: cluster-test Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 1,2,3
四、补充说明
- 脚本废弃原因:Kafka 3.0+ 把「优先副本选举」「手动Leader选举」等功能整合到
kafka-leader-election.sh中,统一了 Leader 操作入口; - 无需JSON文件:新版本命令支持直接指定 topic/partition,无需手动创建JSON配置,更简洁;
- 生产环境建议:非必要不手动切换 Leader(Kafka 懒切换策略更优),仅在需要负载均衡时执行。

综上,你的 Kafka 集群高可用配置正常,故障转移机制生效,集群状态符合生产环境的高可用要求!
Kafka启动异常修复脚本clean_kafka_data.sh
# 在每个 Kafka 节点上执行(kafka1, kafka2, kafka3)
cat > clean_kafka_data.sh << 'EOF'
#!/bin/bash
echo "=== 清理 Kafka 数据并重启 ==="
# 1. 停止 Kafka
echo "停止 Kafka..."
/opt/kafka/bin/kafka-server-stop.sh 2>/dev/null
sleep 3
pkill -9 -f Kafka 2>/dev/null
# 2. 备份并清理 meta.properties
echo "清理 meta.properties..."
if [ -f "/opt/kafka/kafka-logs/meta.properties" ]; then
backup_file="/opt/kafka/kafka-logs/meta.properties.bak_$(date +%Y%m%d_%H%M%S)"
mv /opt/kafka/kafka-logs/meta.properties "$backup_file"
echo "已备份到: $backup_file"
else
echo "meta.properties 文件不存在"
fi
# 3. 启动 Kafka
echo "启动 Kafka..."
/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties
# 4. 等待并检查
echo "等待启动..."
sleep 15
if jps | grep -q Kafka; then
echo "Kafka 启动成功"
else
echo "Kafka 启动失败,请检查日志"
tail -10 /opt/kafka/logs/kafkaServer.out 2>/dev/null
fi
# 5. 查看 Kafka 默认运行端口 9092
netstat -pantul|grep 9092
EOF
chmod +x clean_kafka_data.sh
./clean_kafka_data.sh
重置进度 + 开始消费
kafka1
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server 192.168.88.101:9092 \
--group test-group \
--topic test_topic \
--reset-offsets --to-earliest --execute
kafka2
/opt/kafka/bin/kafka-console-consumer.sh \
--topic test_topic \
--bootstrap-server 192.168.88.102:9092 \
--group test-group
常见报错
1、myid 文件缺失

2、配置报错

单词拼接错误,应该是kafka。
小结
常用命令
# 创建 Topic
/opt/kafka/bin/kafka-topics.sh --create --topic my-topic-1 --bootstrap-server kafka1:9092
# 查看 Topic
/opt/kafka/bin/kafka-topics.sh --list --bootstrap-server kafka1:9092
# 查看 Topic 详情
/opt/kafka/bin/kafka-topics.sh --describe --bootstrap-server kafka1:9092

日志分析
# 查看 Kafka 日志
tail -f /opt/kafka/logs/server.log
# 查看 ZooKeeper 日志
tail -f /opt/kafka/logs/zookeeper.out
# 查看 Zookeeper GC 日志
tail -f /opt/kafka/logs/zookeeper-gc.log
Kafka 架构与工作流程(了解)
技术架构



工作流程
- 生产者定期向主题发送消息。
- Kafka 代理存储为该特定主题配置的分区中的所有消息。 它确保消息在分区之间平等共享。 如果生产者发送两个消息并且有两个分区,Kafka 将在第一分区中存储一个消息,在第二分区中存储第二消息。
- 消费者订阅特定主题。
- 一旦消费者订阅主题,Kafka 将向消费者提供主题的当前偏移(offset),并且还将偏移保存在 Zookeeper 。
- 消费者将定期请求 Kafka (如100 Ms)新消息。
- 一旦 Kafka 收到来自生产者的消息,它将这些消息转发给消费者。
- 消费者将收到消息并进行处理。
- 一旦消息被处理,消费者将向 Kafka 代理发送确认。
- 一旦 Kafka 收到确认,它将偏移更改为新值,并在 Zookeeper 中更新它。 由于偏移在 Zookeeper 中维护,消费者可以正确地读取下一封邮件。
- 以上流程将重复,直到消费者停止请求。
消费者可以随时回退/跳到所需的主题偏移量,并阅读所有后续消息
Topic的分片和副本
分片有什么用呢?
分片是什么?
分片就是分区,就是将数据拆分到不同的节点(服务器)去工作
解决单台节点容量有限的问题, 当然节点(服务器)变多了, 自然整个效率也会提升, 吞吐量也会提升
# 通俗解释:一台服务器的资源(内存,硬盘,cpu等)有限,我用多几台服务器,资源进行整合后,就效率提升了
通过分片, 可以将一个大的容器拆解为多个小的容器 分布在不同的节点上, 从而实现, 让多个服务器来一起存储数据, 实现分布式存储
# 通俗解释:有一个10T的数据,由原来的1台服务器存储,现在交由10台服务器去存储,此时变成 1台服务器存储1T数据
分片的数量有没有限制? 理论上没有限制 , 与节点的数量是没有关系的 分片数量不会超过总节点数量的3倍
分片的数量建议是 节点(kafka启动后每台服务器只有一个节点)的 2-3倍
分片的数据: --partition 服务器数量的2-3倍
副本有什么用呢?
提高数据的可靠性, 副本越多, 数据越可靠, 但是数据冗余程度越高,对于硬盘资源会比较浪费
hadoop就是这样的:
比如:有10台服务器,建议是3个副本,但是,能不能将数据压缩一下
hdfs存储数据的时候,会选择使用压缩来进行数据存储
副本数量是否有限制呢? 理论上有限制, 最多和节点的数量是相等的, 但是一般建议是1~3个副本之间
当副本数是1的时候,就代表,没有副本
Kafka 核心知识与实操总结
一、核心概念
1. MQ 消息队列
-
本质:系统间消息传递的"快递站",缓冲高并发流量,确保消息不丢失、不重复
-
两种模式:
-
- 点对点:消息仅被一个消费者消费,消费后删除
- 发布-订阅:消息存储在主题,多个订阅者可接收所有消息(Kafka采用此模式)
2. Kafka 核心组件
-
三大主体:
-
- Broker:Kafka集群节点,负责存储和处理消息
- Producer:消息生产者,向主题发送消息(可指定分区)
- Consumer:消息消费者,从主题读取消息(通过消费者组实现并行处理)
-
三大内容:
-
- Topic:消息分类单元(类似数据库表),是消息收发的核心载体
- Partition:主题的分片,分散存储在不同Broker,实现并行读写和扩容(单个分区内消息有序)
- Replica:分区的备份,分为Leader(对外提供服务)和Follower(同步数据),保障高可用
-
依赖组件:Zookeeper(传统架构),负责集群协调、Leader选举、元数据存储
3. 关键特性
- 高吞吐:单机支持每秒10万+消息传输
- 持久化:以O(1)时间复杂度存储TB级数据
- 高可用:通过副本机制实现故障自动转移
- 水平扩展:支持Broker、Partition动态扩容
- 消息语义:支持At most once(最多一次)、At least once(最少一次,常用)、Exactly once(仅一次)
二、核心作用
- 限流削峰:缓冲高并发请求,避免后端服务崩溃
- 解耦系统:降低生产者和消费者的直接依赖
- 数据同步:支持离线处理(如Hadoop)和实时处理(如Spark)
- 日志收集:统一收集分布式系统日志,便于分析
- 活动跟踪:存储用户行为数据,支撑推荐、统计等场景
三、部署与实操
1. 环境要求
- 操作系统:CentOS Stream 9(或兼容Linux系统)
- Java版本:Kafka 3.0+推荐Java 11/17
- 端口:Zookeeper(2181)、Kafka(9092)、集群同步(2888/3888)
2. 单机部署关键步骤
- 安装Java并配置环境
- 下载解压Kafka
- 启动内置Zookeeper:
bin/zookeeper-server-start.sh config/zookeeper.properties - 启动Kafka Broker:
bin/kafka-server-start.sh config/server.properties - 验证:
jps查看QuorumPeerMain(Zookeeper)和Kafka进程
3. 集群部署核心要点
-
节点配置:3台服务器(奇数节点便于Zookeeper选举),配置主机名、免密登录、时间同步
-
Zookeeper集群:
-
- 配置
zookeeper.properties,指定集群节点(server.1=IP:2888:3888) - 每个节点创建
myid文件(值为1/2/3,与配置对应)
- 配置
-
Kafka集群:
-
- 配置
server.properties,指定broker.id(唯一)、listeners(节点地址)、zookeeper.connect(集群地址) - 关键参数:
num.partitions=3(默认分区数)、default.replication.factor=3(默认副本数)
- 配置
-
启动顺序:先启动Zookeeper集群,再启动Kafka集群
4. 常用命令(速查)
| 功能 | 命令格式 |
|---|---|
| 创建主题 | kafka-topics.sh --create --topic 名称 --bootstrap-server 地址 --partitions 数 --replication-factor 数 |
| 查看主题列表 | kafka-topics.sh --list --bootstrap-server 地址 |
| 查看主题详情 | kafka-topics.sh --describe --topic 名称 --bootstrap-server 地址 |
| 发送消息(生产者) | kafka-console-producer.sh --topic 名称 --bootstrap-server 地址 |
| 消费消息(历史) | kafka-console-consumer.sh --topic 名称 --bootstrap-server 地址 --from-beginning |
| 消费消息(实时) | kafka-console-consumer.sh --topic 名称 --bootstrap-server 地址 |
| 停止服务 | 先停Kafka:kafka-server-stop.sh,再停Zookeeper:zookeeper-server-stop.sh |
5. 集群验证
- 高可用测试:关闭某Broker,查看主题Leader是否自动切换,消息仍可正常收发
- 数据同步:多节点生产/消费消息,验证所有消费者接收数据一致
- 副本状态:
describe topic查看Isr(同步副本列表)包含所有节点
四、核心设计理解
-
分区(Partition)作用:
-
- 突破单节点存储限制,实现分布式存储
- 并行读写,提升吞吐量(多个分区可同时被消费)
-
副本(Replica)作用:
-
- 数据冗余,防止节点故障导致数据丢失
- Leader故障时,Follower自动选举为新Leader,保障服务不中断
-
数据存储:默认存储在
/tmp/kafka-logs(建议修改为持久化目录如/var/lib/kafka/data)