
前一阵我在本地想验证一个消费组重平衡的逻辑随手起了Kafka默认的单实例结果发现主题副本永远只有1个想测故障迁移只能手动改配置折腾了半天差点放弃。后来我换了个思路单机多broker进程也就是大家常说的伪分布式Kafka。在同一台机器上跑三个broker端口和数据目录分开整个集群的多数核心机制都能正常测。这篇赫兹威客系列的测试教程就把我的搭建流程和踩坑记录完整写下来适合想最快搞定Kafka测试环境的人也适合准备Kafka面试题时需要真实环境验证结论的读者。所谓伪分布式本质上是模拟而不是模拟它用多个进程替代了多台机器让你在本地就能体验副本复制、Leader选举、ISR同步这些分布式特性。这样搭出来的环境虽然不能用于性能压测但对于学习原理、调试客户端代码、复现线上问题来说性价比极高。下面我就从设计思路开始一步一步拆解整套搭建过程。1. 为什么测试环境要搭伪分布式Kafka集群1.1 伪分布式到底解决了什么问题很多同学学Kafka的第一步是下载官方压缩包解压后直接跑默认配置。这样确实能启动一个Kafka进程也能发送和接收消息但一旦你需要测副本机制、消费组重平衡、Broker故障恢复这些核心功能单实例就会立刻露馅。单实例下你无法验证以下场景副本因子设置大于1时主题创建会直接失败或警告一个Broker宕机后分区Leader能否自动切换到其他副本ISRIn-Sync Replicas列表如何收缩和恢复同组多个消费者如何分配分区重平衡如何触发__consumer_offsets 主题的多副本机制是否正常伪分布式Kafka破解了这个问题。在单台机器上启动多个Broker进程让它们互相认识、组成一个逻辑上的集群。这样上面提到的分布式行为都可以被真实触发而代价仅仅是一台机器的内存和磁盘。1.2 伪分布式和真集群的差别以及什么时候用它伪分布式不是Kafka官方术语它借用了Hadoop生态里的叫法。在Hadoop里伪分布式是指用多个进程模拟HDFS和YARN的多节点部署在Kafka里意思类似指多个Broker进程跑在同一台机器上通过不同端口区分。和真正的分布式集群相比它们没有架构上的本质区别因为Kafka的Broker本来就是一个独立进程彼此通过TCP通信。区别主要在于部署位置和资源隔离程度。我整理了一张对比表方便你判断对比项单实例伪分布式多机真集群进程数1多个多个部署范围单机单机多台机器副本机制无法验证可完整验证可完整验证故障演练没有意义可手动kill进程可物理级演练性能参考完全不能参考基本不能参考可做基准压测运维复杂度低中高所以我的建议是如果你只是写API Demo单实例完全够用如果你想弄懂Kafka的分布式原理或者要写一套生产级客户端代码伪分布式是最合适的实验环境而只有当你准备上生产才需要考虑多机真集群。这套环境还有一个额外价值——面试前突击。很多Kafka面试题都会问“副本因子和ISR的关系”“Leader选举是怎么触发的”“消费组重平衡发生什么”这些问题背答案容易忘但如果你在伪分布式环境里亲手杀过一次Broker看过一次分区状态变化再去面试心里就有底了。2. 环境准备JDK、ZooKeeper、Kafka 版本选型2.1 版本组合怎么选Kafka组件对版本兼容比较敏感特别是JDK和ZooKeeper。我这次用的是当前比较主流的一套组合JDK 11 ZooKeeper 3.8.1 Kafka 3.4.1。其中Kafka 3.x版本同时支持JDK 8和JDK 11如果你的机器只有JDK 8也不影响操作但建议优先用JDK 11多版本管理工具如SDKMAN或jenv会方便不少。这里要提醒一句如果你用的是Kafka 3.0以上版本ZooKeeper模式依然是默认路径但Kafka已经引入了KRaft模式可以不依赖ZooKeeper直接跑。我这次教程依然采用ZooKeeper模式因为存量系统里它依然是主流而且理解ZooKeeper的功能有助于你理解Broker注册、Controller选举这些机制。想尝试KRaft的话可以把同一套配置思路平移过去只是把ZooKeeper相关配置替换为controller.quorum.voters。单机伪分布式环境对硬件要求不高但三个Broker进程同时跑建议内存至少4GB磁盘留出5GB以上。如果机器配置太低可以只启动两个Broker后续所有操作也够用只要别在创建主题时把副本因子设到3。2.2 下载与基础配置Kafka和ZooKeeper都建议从Apache官网下载正式发布版本下载地址是kafka.apache.org/downloads如果网络慢可以用国内镜像源比如清华源或阿里源。Kafka的二进制包是tgz格式直接解压即可不需要编译。解压后我把目录整理成下面这样方便后面统一管理~/kafka-lab/ ├── kafka_2.13-3.4.1/ # Kafka主目录 │ ├── bin/ │ ├── config/ │ ├── libs/ │ └── log4j.properties ├── zookeeper-3.8.1/ # ZooKeeper主目录 │ ├── bin/ │ ├── conf/ │ └── lib/ ├── data/ │ ├── zk-data/ # ZooKeeper数据目录 │ ├── kafka-logs-1/ # Broker 1数据目录 │ ├── kafka-logs-2/ # Broker 2数据目录 │ └── kafka-logs-3/ # Broker 3数据目录Windows用户需要注意Kafka和ZooKeeper的启动脚本分别是bin目录下的.sh和.bat版本Windows下用.bat即可但目录路径不要带空格否则脚本容易解析出错。JDK环境变量JAVA_HOME必须提前配好否则启动时会直接报“找不到Java”的错误。解压完成后给三个数据目录赋好权限确保当前用户可读写。然后检查一下Java版本java -version如果输出显示openjdk version 11.0.x环境就绪。这个阶段还有一个细节容易被忽略ZooKeeper默认会占用2181端口三个broker会分别占用9092、9093、9094端口本机没有其他服务占用这些端口即可。3. 单机多broker核心搭建同一台机器模拟三节点集群3.1 broker配置文件差异化要点这是整个伪分布式搭建中最核心的部分。Kafka的配置主要在config/server.properties单实例部署直接用它就行但我们要启动三个broker就需要准备三个独立的配置文件。每个文件里必须差异化配置四个核心项broker.id、listeners、log.dirs、以及日志相关路径。先看第一个broker的配置文件。我复制server.properties为server-1.properties修改关键内容如下# config/server-1.properties broker.id1 listenersPLAINTEXT://localhost:9092 advertised.listenersPLAINTEXT://localhost:9092 log.dirs/home/test/kafka-lab/data/kafka-logs-1 zookeeper.connectlocalhost:2181 offsets.topic.replication.factor3 transaction.state.log.replication.factor3 transaction.state.log.min.isr2 auto.create.topics.enablefalse第二个和第三个broker配置类似只需要修改差异项# config/server-2.properties broker.id2 listenersPLAINTEXT://localhost:9093 advertised.listenersPLAINTEXT://localhost:9093 log.dirs/home/test/kafka-lab/data/kafka-logs-2 # config/server-3.properties broker.id3 listenersPLAINTEXT://localhost:9094 advertised.listenersPLAINTEXT://localhost:9094 log.dirs/home/test/kafka-lab/data/kafka-logs-3zookeeper.connect保持一致指向同一个ZooKeeper节点。这样三个broker启动后会在ZooKeeper上注册到同一个集群互相同步元数据。为了更直观地看出差异我列了一个表格配置项broker 1broker 2broker 3broker.id123listenerslocalhost:9092localhost:9093localhost:9094advertised.listenerslocalhost:9092localhost:9093localhost:9094log.dirskafka-logs-1kafka-logs-2kafka-logs-3这里有几个点你必须理解否则后面会踩坑。第一broker.id是集群内唯一标识用于区分不同broker不能重复。第二listeners是broker对外提供服务的地址和端口同一台机器上必须用不同端口避免冲突。第三advertised.listeners是broker注册到ZooKeeper后对外公布给客户端和其他broker的连接地址。如果配置的是localhost那么只有本机客户端能连如果其他机器要访问这里必须改成机器的实际IP否则会出现“能连上9092端口但客户端却反复拉取不到元数据”的诡异问题。第四log.dirs是消息数据存储目录每个broker必须使用独立目录共享目录会导致数据错乱甚至broker启动失败。第五offsets.topic.replication.factor设置为3是为了让内部消费组位移主题也具备三副本。在单broker环境下这个值必须改成1否则消费组功能会报错在三个broker的伪分布式环境中保留3就是合理的。3.2 启动顺序与初始化验证Kafka集群启动顺序有讲究先启动ZooKeeper再启动各个Kafka broker。如果ZooKeeper没有起来broker启动会一直重试连接并报错。启动ZooKeeper我用独立ZooKeeper发行版cd ~/kafka-lab/zookeeper-3.8.1 bin/zkServer.sh start如果你用的是Kafka自带的ZooKeeper脚本也可以cd ~/kafka-lab/kafka_2.13-3.4.1 bin/zookeeper-server-start.sh config/zookeeper.properties但注意这种方式的zookeeper.properties路径要改成你自己配置的ZooKeeper数据目录。然后启动三个brokercd ~/kafka-lab/kafka_2.13-3.4.1 export KAFKA_HEAP_OPTS-Xmx512M -Xms256M bin/kafka-server-start.sh config/server-1.properties bin/kafka-server-start.sh config/server-2.properties bin/kafka-server-start.sh config/server-3.properties 我特意把KAFKA_HEAP_OPTS调低了默认堆内存比较大三个broker同时跑容易把机器内存吃满。每个broker给512MB堆内存足够完成测试操作。启动后用jps命令查看进程jps -l正常情况下会看到三个Kafka进程和一个ZooKeeper进程一共4个Java进程。然后去看每个broker的日志文件确认没有报错。日志默认打印在启动命令的终端也可以在log.dirs类似路径下查看server.log。看到“started (kafka.server.KafkaServer)”字样就说明broker启动成功了。3.3 验证集群拓扑与元数据进程启动不代表集群状态正常。在ZooKeeper上检查broker是否全部注册是最直接的验证方式。使用ZooKeeper命令行客户端cd ~/kafka-lab/zookeeper-3.8.1 bin/zkCli.sh -server localhost:2181 ls /brokers/ids正常会输出[1, 2, 3]表示三个broker都注册到了同一个集群。如果只看到一个或两个ID说明有的broker没连上ZooKeeper回去看对应日志。还可以用Kafka自带的命令行查看broker版本信息cd ~/kafka-lab/kafka_2.13-3.4.1 bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092这个命令能列出集群中所有broker的ID和API版本信息结果里能看到3个broker就说明拓扑正常。到这里伪分布式集群已经搭建完成接下来就可以进行各种功能验证了。4. 生产消费测试与故障演练4.1 创建带副本的主题集群搭好之后第一件事就是创建一个多副本主题验证副本机制是否真正生效。我以一个名为test-replica的主题为例3个分区、3个副本cd ~/kafka-lab/kafka_2.13-3.4.1 bin/kafka-topics.sh --bootstrap-server localhost:9092 \ --create --topic test-replica \ --partitions 3 --replication-factor 3创建后查看主题状态bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic test-replica输出大概长这样Topic: test-replica TopicId: xxxx PartitionCount: 3 ReplicationFactor: 3 Topic: test-replica Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3 Topic: test-replica Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1 Topic: test-replica Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2Replicas一列显示的是该分区副本分布在哪些broker上Isr是当前处于同步状态的副本。三列都是3个broker说明多副本创建成功了。如果你在实际操作中发现Replicas始终只有1个最常见的原因有两个一是broker没全部启动二是创建主题时没有指定replication-factor参数使用了默认值1。伪分布式测试中请在创建主题时显式指定副本数。这里要注意offsets.topic.replication.factor3意味着内部消费组位移主题也要求3个副本。如果某个broker没起来创建消费组时可能报“Error while fetching offset topic”之类的错误。这也是为什么我建议三个broker都要启动正常再开始测试。4.2 生产消费测试与指定时间消费主题创建好了用命令行生产者发几条消息验证链路bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test-replica输入几条消息比如hello kafka、pseudo distributed test然后按CtrlC退出。再启动消费者bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic test-replica --from-beginning如果能输出刚才输入的消息生产消费链路就是通的。这里有个新手容易踩的坑消费者不加--from-beginning时默认只消费启动之后新到达的消息不读取历史数据所以看起来就像“收不到消息”。热词里有人问“kafka消费命令指定消费时间”这在实际排查中很常用。伪分布式集群里也能验证。如果你想从某个时间点开始消费可以这样重置消费组的offsetbin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group test-group --topic test-replica \ --reset-offsets --to-datetime 2024-01-01T00:00:00.000 --execute执行前需要确保test-group消费组已经存在。如果只想消费某个分区的特定offset范围用console-consumer直接指定分区和offsetbin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic test-replica --partition 0 --offset 100 --max-messages 10这条命令会从分区0的第100条消息开始最多消费10条。对调试消息积压、验证某个时间点后的数据非常实用。4.3 故障演练杀掉一个broker看集群表现伪分布式环境最大的价值就是可以大胆做故障演练。我现在模拟broker 1宕机按ctrlc终止第一个broker进程或者直接用kill命令杀掉PID。记住这里要kill的是kafka-server进程不是ZooKeeper。杀掉后再次查看主题状态bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic test-replica你会发现原本由broker 1担任leader的分区leader自动切换到了broker 2或broker 3。比如原先Partition 0的Leader是1现在变成了2。同时ISR列表中broker 1消失了只剩下2和3。这就是Kafka的高可用机制在起作用只要还有同步副本存活分区就能继续提供服务消息不会丢失。然后模拟broker 1恢复重新启动它等它重新加入集群。再次查看describe你会发现broker 1会重新出现在Isr列表中分区副本也恢复到3个。这个“杀了又活”的过程就是面试题里常说的Leader选举和ISR收缩恢复。我还推荐做一个更极端的验证在生产消息的同时杀掉breeder观察生产者是否有报错。如果acksall且主题的min.insync.replicas配置合理短暂故障期间生产可能报超时但恢复后能继续生产。这能帮你理解acks参数和可用性之间的权衡。顺带说一句热词里“kafka消息延迟高”也是大家关心的重点。在伪分布式环境里排查延迟可以先看消费组Lagbin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --describe --group test-group输出中LAG列表示消费者落后多少条消息。如果LAG持续增长通常要么是消费者处理太慢要么是分区数不足导致并发度上不去。在伪分布式环境中单机资源有限不要拿它做高并发压测但用这个环境验证消费组重平衡、验证分区分配逻辑是完全够用的。5. 可视化工具与常用辅助排查手段5.1 图形化工具选型命令行用多了总觉得没有一个直观的界面很痛苦。好在Kafka生态里可视化工具不少我推荐三款按场景选择。第一款是Offset Explorer以前叫Kafka Tool桌面客户端Windows和macOS都能装。它支持多集群配置可以查看主题列表、分区详情、消息内容、消费组offset适合日常快速浏览。连接配置时填上bootstrap服务器地址比如localhost:9092就可以直接连上我们的伪分布式集群。第二款是Kafka UI开源Web界面支持多集群管理界面现代可以查看主题、消费者组、消息、还支持发送测试消息。部署方式可以跑在Docker里也可以直接用Java进程启动配置文件中指定kafka集群地址即可。适合团队内部搭一个共享的Kafka管理界面。第三款是Kafdrop非常轻量主打消息浏览和主题查看启动参数简单适合临时使用。我自己在伪分布式环境中用得最多的是Offset Explorer因为桌面端部署简单点开就能看到三个broker是否都在线也很直观地展示分区的leader和副本分布。注意这些工具连接多broker集群时填其中一个broker地址即可Kafka客户端会自动拉取整个集群的元数据。5.2 几种必会的命令行排查手段图形化工具适合浏览但真要定位问题命令行依然是最高效的。除了前面用过的kafka-topics.sh和kafka-consumer-groups.sh还有一个必会工具是查看某个主题某个分区的起始和最新offsetbin/kafka-get-offsets.sh --bootstrap-server localhost:9092 \ --topic test-replica --time -1--time -1表示查最新offset--time -2表示查最早offset。这个命令配合consumer-groups的LAG信息能快速判断消费位点是否已经过期或丢失。如果遇到broker启动异常第一件事是看日志。Kafka的日志文件路径由log.dirs指定每个broker目录下会有server.log。日志通常已经把问题原因写得比较清楚比如端口占用、配置错误、权限不足等。不要一上来就怀疑玄学先看日志再找配置。另外可以开启JMX监控。Kafka本身支持JMX启动前设置JMX端口即可export JMX_PORT9999然后用JConsole连接本地端口或配合Prometheus的JMX exporter采集指标。这套东西在伪分布式环境就能跑通等以后上生产可以无缝平移监控方案。热词里有“kafka exporter下载”说的就是这个JMX exporter网上可以找到现成的jar包和配置文件。6. 常见问题排查与避坑清单6.1 高频问题速查表我在搭建和测试过程中遇到过不少问题有些是低级的配置错误有些是Kafka机制导致的迷惑行为。整理成一张速查表按症状、原因、解决方式排列你可以直接对照排查。症状常见原因解决办法broker启动失败提示端口已被占用9092/9093/9094端口被其他进程占用lsof -i:9092 查看占用进程关闭或换端口启动后jps看不到3个broker内存不足或堆内存配置过大设置KAFKA_HEAP_OPTS为512M确保本机至少有4G内存创建主题时副本因子始终为1创建命令没加--replication-factor参数创建时显式指定--replication-factor 3消费组创建报错提示offsets topic有问题__consumer_offsets主题副本数不满足确认offsets.topic.replication.factor3且三个broker都在线客户端能连9092但拉取元数据失败advertised.listeners设置成localhost而客户端从其他机器访问listeners和advertised.listeners都改成实际IP消费者收不到历史消息没有加--from-beginning加参数或按时间重置offset杀掉一个broker后某个分区不可用该分区的ISR中只剩leader且min.insync.replicas设置过高等待broker恢复或调低min.insync.replicasWindows下启动脚本报错路径带空格或JAVA_HOME未配置把Kafka放在无空格目录检查JAVA_HOME消息延迟高LAG不断增长分区数少于消费者数部分消费者空转增加分区数或调整消费者并发逻辑6.2 我踩过几次坑之后的心得这套伪分布式环境我搭过不止一次有几点心得想特别分享。第一配置文件的独立性比想象中更重要。很多人图省事想用一个配置文件启动三个broker进程结果各种数据目录冲突、broker.id重复根本起不来。配置文件复制三份每个改动三到四个参数是最稳妥的做法。宁可多花两分钟写配置文件也不要省这一步。第二测试完一定要清理数据目录。Kafka和ZooKeeper会把元数据和消息数据写到磁盘如果你反复重建集群旧数据没有清理可能触发各种奇怪问题比如分区元数据对不上、offset错乱。重新搭建测试环境时把data目录下的zk-data和kafka-logs-*全部删掉再启动一般都能恢复正常。第三伪分布式环境的性能数据真不能当参考。三个broker共享同一块磁盘、同一块网卡、同一套CPU它们的“三节点”性能甚至可能不如一个独立的物理节点。我见过有人拿这种环境测出“高吞吐”数据后直接用于生产容量规划这个误区很危险。伪分布式适合验证功能正确性和机制逻辑不适合做任何性能基准。最后再说一个实用小技巧。如果你在测试消费组重平衡可以打开控制台消费者日志的DEBUG级别在里面能看到详细的JoinGroup、SyncGroup、Revoke分区记录。开启方式是在log4j.properties里把kafka.coordinator.group的日志级别改成DEBUG然后重启消费者进程。这样你能直观看到分区rebalance的完整过程对理解Kafka消费组原理帮助很大。