ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka搭建实战:从单机到Docker及常见问题排查

Kafka搭建实战:从单机到Docker及常见问题排查 后端开发做到一定阶段消息队列基本是绕不开的。我在好几个项目里都被问到同一个问题如果要做高并发的系统队列选谁答案很多时候是Kafka。这个系列我打算从搭建开始一步步把Kafka从部署讲到生产消费、再到运维排查。第一篇文章就先把服务搭起来顺便把那些搭建时最容易踩的坑一次性讲清楚。先说清楚这篇内容适合谁刚接触消息队列、想在本机搭一套Kafka做实验的同学团队里要快速搭建开发/测试环境、但不想再被各种配置折腾的后端以及已经把Kafka跑起来但遇到连接不上、元数据拉取失败这类问题想找排查思路的朋友。我会尽量用“为什么这么做”的逻辑来讲而不是光贴命令。1. 先聊聊Kafka它解决的到底是什么问题很多人一上来就装环境、跑脚本结果Kafka是装好了但为什么要用它、里面那些概念对应什么场景还是一头雾水。我建议先花十分钟把定位搞清楚后面写代码、调参数时才不会懵。1.1 消息队列的三大作用削峰、解耦、异步面试和实际项目里最常说的就是这三个我按真实场景给你拆开。异步。用户注册成功后系统要发欢迎短信、发优惠券、记录埋点日志。如果这些都在注册接口里同步做每增加一个步骤接口响应就慢一截。引入消息队列后注册服务只负责把“用户注册成功”这个事件写到Kafka其他服务自己去消费这个消息做各自的事。注册接口的耗时从“等所有下游完成”变成“等消息写入完成”体感快很多。削峰。典型的场景是秒杀、抢购。平日里系统流量平稳一到大促瞬间涌入几万甚至几十万的请求直接把数据库打挂。用Kafka挡在前面先把请求消息接住后端服务按自己的消费能力一条一条处理。这就是把突刺流量“削平”了系统不会因为瞬时流量崩溃。解耦。修改一个核心流程最怕的是牵一发动全身。比如订单系统要通知库存、通知积分、通知物流如果都走接口调用任何一个下游改动上游也要跟着改。用消息队列把上下游隔开上游只负责生产消息下游自己订阅互不干扰。新增一个下游消费者上游代码一行都不用动。这三个作用不是独立的实际项目中往往是同时体现。理解了这个你就知道为什么几乎所有中大型后端系统里都有消息队列的身影。1.2 入门必须搞懂的四个概念Topic、Partition、Offset、Consumer Group搭建Kafka之前这几个概念一定要有初步印象因为之后的配置、命令全在跟它们打交道。Topic是消息的逻辑分类相当于数据库里的表。一条订单消息发到order_topic一条日志消息发到log_topic互不干扰。Partition是Topic下面的物理分片。一个Topic可以拆成多个Partition每个Partition是一个有序的消息日志。为什么要分片因为单一个文件写起来有瓶颈拆成多个Partition就可以并发读写还能分布到不同机器上吞吐量才能上去。Kafka只保证Partition内部有序不保证Topic全局有序这个限制很多新手没注意。Offset是消息在Partition里的位置编号类似数组下标。消费者每消费一条消息就把Offset往后挪。如果消费者宕机重启可以从上次记录的Offset继续消费不会丢消息。Consumer Group是消费者的分组机制。同一个Group里的消费者共同消费一个Topic一条消息只被组内某一个消费者处理不同Group之间相互独立消息会被每个Group都消费一遍。我举个好记的例子一条消息就是一份通知一个Group是一家公司。通知发到公司AA内部谁有空谁处理处理一次就行公司B也能收到同一份通知但A和B的处理互不影响。这四个概念搭建阶段你只需要眼熟它们真正写生产消费代码的时候会反复用到。1.3 选型对比Kafka、RabbitMQ、RocketMQ我为什么先学它很多新手选型时会纠结网上对比文章也很多。我说下我的真实感受维度KafkaRabbitMQRocketMQ吞吐量极高百万级每秒中等万级每秒高十万级每秒消息可靠性可配置默认在性能与可靠之间平衡较高支持多种确认机制较高支持事务消息消息路由以Topic为主功能简单路由规则灵活Exchange以Topic为主支持标签过滤延迟消息/定时消息原生不支持需额外设计支持死信队列配合实现原生支持学习成本概念少命令简单生态庞大概念多但资料也多偏重适合Java技术栈典型场景大数据管道、日志收集、高吞吐事件流企业应用、微服务内部异步电商、金融等对可靠性要求高的业务我的建议很直接如果你想在最短时间内把消息队列跑起来并且后续会接触大数据生态Flink、Spark、ClickHouse这些优先学Kafka。它的设计朴素核心就是“一个超高吞吐的分布式提交日志”没有太多花哨能力正因如此反而容易掌握核心思想。等Kafka用熟了再看RabbitMQ、RocketMQ你会发现很多概念是相通的只是实现方式不同。2. 搭建前的准备工作JDK版本、下载与目录结构搭建Kafka不需要多高深的操作但准备工作做对了后面能省很多事。我见过太多人卡在版本不匹配上所以这一章把选版本、下载、目录结构一次讲明白。2.1 检查JDK版本别选错Kafka是用Java和Scala写的运行必须依赖JVM。现在Kafka 3.x版本要求JDK 8以上官方推荐JDK 11或17。我建议直接用JDK 17理由很简单性能更好而且Kafka 4.x也继续支持不会一升级系统就跟着换JDK。检查命令java -version如果没装或者版本是1.7这种古董级先去装一个JDK 17。装的时候注意Linux下配置好JAVA_HOME环境变量Windows下注意Path里不要混入多个JDK版本这类问题排查起来很费时间。2.2 下载Kafka官方包还是源码编译直接下载官方编译好的二进制包就行千万别自己去源码编译没必要。Apache官网下载慢的话用国内镜像源搜索“Kafka 清华镜像”就能找到版本更新速度也能跟上。选版本有个小技巧不要追最新的大版本选当前稳定版即可比如3.6到3.9之间的版本都行。如果你要看别人写的教程尽量选择与教程相同的大版本比如3.x系列避免因为版本差异导致配置项对不上。Kafka 4.0以后移除了Zookeeper依赖配置方式有变化老教程里的Zookeeper那套在4.0上就不适用了。下载得到的包名类似kafka_2.13-3.7.0.tgz前面的2.13是Scala编译版本后面的3.7.0才是Kafka版本号。解压tar -xzf kafka_2.13-3.7.0.tgz cd kafka_2.13-3.7.02.3 解压后的目录结构bin、config、libs分别干嘛用的很多人解压完就急着启动结果找不到命令先花两分钟认识下目录目录作用bin存放启动脚本和命令行工具里面有.sh和.bat两种Windows用后者config配置文件所在地server.properties、zookeeper.properties都在这里libsKafka依赖的所有Jar包一般不用动logs运行日志目录默认不存在的首次启动后生成这里特别注意config目录。Kafka的很多“诡异问题”都是配置不对导致的尤其是监听地址、Zookeeper连接地址这两个后面实操部分我会单独展开。3. 单机版搭建实操从修改配置到收发第一条消息这一章是全文最核心的实操部分。我会用最经典的模式来做一个Zookeeper 一个Kafka Broker。虽然Kafka现在已经有KRaft模式可以不要Zookeeper但绝大多数线上老项目和教程仍然是Zookeeper模式先把它弄明白后面再切KRaft会非常轻松。3.1 server.properties 核心配置项怎么改进入config目录编辑server.properties。一个单机开发环境真正需要改的不多关键是理解每项的含义。broker.id0 listenersPLAINTEXT://:9092 advertised.listenersPLAINTEXT://localhost:9092 log.dirs/tmp/kafka-logs zookeeper.connectlocalhost:2181 num.partitions3 log.retention.hours168逐项解释broker.idKafka集群中每个节点唯一ID单机就填0。listenersBroker对外监听的地址和端口。默认是9092这行决定Kafka在哪个网卡上接受连接。advertised.listeners这个最容易踩坑。它是Broker告诉客户端“你应该连接我哪个地址”。如果本机访问没问题但局域网内其他机器连不上十有八九就是这项配置不对。填localhost只适合本机测试要让别的机器访问这里要填IP或域名。log.dirs消息数据落盘的目录。开发环境放/tmp可以但生产环境务必放在数据盘并且目录空间要监控。zookeeper.connectZookeeper地址。Kafka依赖Zookeeper保存Broker元数据、Controller选主等。单机默认localhost:2181。num.partitions新建Topic时默认的分区数。如果创建Topic时不指定--partitions就用这个值。开发环境设成3比较合适既能看到分区效果又不会太浪费资源。log.retention.hours消息保留时间默认168小时也就是7天。超过时间的消息会被清理掉。想保留更久就调大想省磁盘就调小。3.2 启动Zookeeper和Kafka并验证进程状态启动顺序有讲究先启动Zookeeper再启动Kafka因为Kafka启动时要向Zookeeper注册信息。# 启动Zookeeper bin/zookeeper-server-start.sh config/zookeeper.properties # 另开一个终端启动Kafka bin/kafka-server-start.sh config/server.properties启动后验证进程是否正常jps正常情况下能看到三个进程QuorumPeerMainZookeeper、Kafka、Jps自身。如果jps看不到Kafka去logs/server.log看报错。如果你用的是新的KRaft模式启动命令不同bin/kafka-storage.sh random-uuid bin/kafka-storage.sh format -t uuid -c config/kraft/server.properties bin/kafka-server-start.sh config/kraft/server.propertiesKRaft模式在下一章用Docker演示这里先不展开。3.3 用命令行脚本创建Topic、生产消息、消费消息服务起来后立刻做一个端到端的验证确保Kafka真的能用。第一步创建Topicbin/kafka-topics.sh --create --topic test-topic \ --bootstrap-server localhost:9092 \ --partitions 3 --replication-factor 1参数含义--topicTopic名称。--bootstrap-serverKafka地址。注意老版本用--zookeeper指定地址现在新版本推荐用--bootstrap-server直接连Broker。--partitions分区数。因为单机只有一个Broker使用must be between 1 and 3的副本数限制所以副本数填1。第二步启动生产者bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test-topic回车后会进入交互模式输入一行就是一个事件。比如输入hello kafka my first message第三步另开终端启动消费者bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic test-topic --from-beginning--from-beginning表示从最早的消息开始消费。如果顺序没错你能在消费者终端看到刚才输入的两条消息。看到消息的那一刻你的第一个Kafka服务就算真正跑通了。验证Topic详情bin/kafka-topics.sh --describe --topic test-topic --bootstrap-server localhost:9092输出里能看到这个Topic有几个分区、每个分区的Leader在哪、副本分布在哪些Broker上。这些信息后面排查消息乱序、堆积问题时很有用。3.4 用Offset Explorer连接本地单机Kafka命令行验证完我建议再装一个可视化工具日常看Topic、看Offset、查消息会方便很多。常用的有Offset Explorer原Kafka Tool、Kafka UI、Kafka Map。我最常用的是Offset Explorer界面直观连接配置也简单。下载安装后界面里添加Cluster点击Add Cluster。Cluster Name随意填比如Local。在Bootstrap Servers输入localhost:9092。如果需要安全认证填对应的Security设置本地单机不需要。点Test Connection显示成功就点Add。连接成功后你能在左侧看到所有Topic列表点进某个Topic能看到Partition分布、Offset范围、消息内容。很多人在这一步遇到“连接成功但看不到Topic”的情况原因是Kafka默认开了auto.create.topics.enable但Topic要等有消息写入后才真正创建空连接时列表就是空的。你只要先启动生产者发送一条消息再用工具刷新就能看到了。4. Docker快速搭建KRaft模式彻底告别Zookeeper如果你不想在本地装一堆依赖或者想快速搭建一套干净的测试环境Docker是更好的选择。这一章用KRaft模式来搭——这是Kafka新推荐的部署方式不再需要Zookeeper。Kafka 4.0已经强制走这套了所以新项目直接学它更省事。4.1 KRaft是什么为什么推荐新项目直接用KRaft是Kafka自己实现的共识协议用Raft算法让Kafka节点自己管理元数据替代外部依赖Zookeeper。以前Kafka要组件再搭一套Zookeeper集群节点多、运维重、脑裂风险也多。KRaft模式下一个Kafka节点可以同时担任两种角色Controller管理元数据、选主和Broker处理消息读写。好处很明显部署组件少一个容器就能跑。没有Zookeeper集群需要维护。元数据同步效率更高集群扩展时更稳。对于开发环境用KRaft模式几秒钟就能起一个Kafka这对需要频繁创建销毁环境的场景太友好了。4.2 用Docker Compose一键启动KRaft版Kafka下面这个Compose配置来自我项目的开发环境直接在干净机器上测试过。用bitnami/kafka镜像环境变量比较多我都加了注释version: 3.8 services: kafka: image: bitnami/kafka:3.7 container_name: kafka-kraft ports: - 9092:9092 environment: # 节点ID单节点固定0 - KAFKA_CFG_NODE_ID0 # 进程角色broker和controller都由这一个节点承担 - KAFKA_CFG_PROCESS_ROLEScontroller,broker # 集群的controller节点列表单节点就是自己 - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS0kafka:9093 # 监听器PLAINTEXT用于客户端通信CONTROLLER用于节点内部通信 - KAFKA_CFG_LISTENERSPLAINTEXT://:9092,CONTROLLER://:9093 # 对外宣告的地址客户端连的就是这个 - KAFKA_CFG_ADVERTISED_LISTENERSPLAINTEXT://localhost:9092 # 监听器协议映射 - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAPCONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT # 指定Controller监听器名称 - KAFKA_CFG_CONTROLLER_LISTENER_NAMESCONTROLLER # 允许自动创建Topic开发环境方便 - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLEtrue volumes: - kafka_data:/bitnami/kafka volumes: kafka_data:启动命令docker-compose up -d等几秒钟确认状态docker-compose ps看到kafka-kraft状态是Up基本就成功了。想确认日志没报错用docker logs kafka-kraft --tail 1004.3 容器内的验证操作和端口映射说明在容器内做一次生产消费验证。先进入容器docker exec -it kafka-kraft bash然后执行# 创建Topic kafka-topics.sh --create --topic docker-test \ --bootstrap-server localhost:9092 \ --partitions 1 --replication-factor 1 # 启动生产者先运行输入消息后按CtrlC退出 kafka-console-producer.sh --broker-list localhost:9092 --topic docker-test # 新建一个终端进入容器启动消费者 docker exec -it kafka-kraft bash kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic docker-test --from-beginning这里要特别留意advertised.listeners的作用。容器里的Kafka进程不知道自己被映射到了宿主机的9092端口它告诉客户端的地址由KAFKA_CFG_ADVERTISED_LISTENERS决定。如果你从宿主机运行Java代码连接Kafka配置里写的地址必须和这个PLAINTEXT://localhost:9092保持一致。如果客户端跑在另一台机器上这里要写宿主机IP比如PLAINTEXT://192.168.1.10:9092。很多人Docker启动后本机客户端连不上基本都是这个地址写错了或者容器映射端口只映射了9092没映射9093。容器间通信、跨机器通信时这部分是排查重点。5. 搭建阶段的高频问题与排查记录最后这一章我把自己和身边同事实际遇到过的问题整理成一份速查表。这些问题在搭建阶段特别容易出现而且报错信息往往比较抽象第一次遇到会卡很久。5.1 连不上Brokererror while fetching metadata with correlation id这是Kafka新手会遇到的最经典错误没有之一。报错长这样Error while fetching metadata with correlation id 3 : {test-topicLEADER_NOT_AVAILABLE}或者客户端直接报超时。看到这个错误按顺序排查Broker是否真的启动了先jps或者docker ps确认进程在。端口是否可通在客户端机器上执行telnet localhost 9092如果不通检查防火墙或安全组。advertised.listeners是否正确这是最隐蔽的坑。客户端能连上Broker但Broker返回给自己的元数据地址是错的。比如客户端在宿主机Broker告诉它“去连kafka:9092”而宿主机根本不认识kafka这个主机名就会报元数据错误。Zookeeper是否正常经典模式下Zookeeper没起来Broker虽然能启动但元数据不完整也会导致这个错误。LEADER_NOT_AVAILABLE这个子类型通常是因为分区Leader还没选举完成多等几秒重试就好。如果一直这样看Broker日志确认分区副本状态常见原因是对应分区的副本不足比如副本因子写了2但只有1个Broker。5.2 消费者重复消费到底是谁的锅这个问题的经典场景是消费者处理完消息但结果又被消费了一遍。很多人第一反应是“Kafka丢消息了”——其实Kafka在这里的表现不是丢而是重复投递。原因通常有几个Offset提交失败。消费者每消费完一批消息要向Broker提交一次Offset。如果消息处理耗时超过max.poll.interval.ms默认5分钟Consumer会触发Rebalance。Rebalance后新的消费者再从上次提交的Offset继续消费上次没提交的那些消息自然被重新拉了一遍。处理完消息但没等提交就宕机。消息处理成功了但Offset还没提交。重启后从旧Offset开始消费这条消息就被重复处理。规避方式消费逻辑要做到幂等重复处理不影响结果。把enable.auto.commitfalse改为手动提交Offset在消息真正处理完之后再提交。这是生产环境的标配做法。调大max.poll.records或减少单批处理量避免一次拉取太多消息导致处理超时。这个问题的根子多在消费者自身不是Kafka服务端的问题排查方向要对。5.3 Kafka消息延迟高和搭建相关的几个原因有时候消息生产了却迟迟到不了消费者。除了消费者本身处理慢搭建阶段导致的几个原因也很常见分区数不合理。Topic分区数太少消费者并发上不去消费速率自然低。生产环境一个Topic的分区数建议根据消息量和消费端并行度提前规划要在创建Topic时定好。Kafka支持增加分区但不建议频繁调整。磁盘性能差。Kafka的强项是顺序写磁盘如果落盘目录在机械硬盘或者IOPS很低的云盘上性能会被拖垮。生产环境用SSD并且log.dirs配置在多块数据盘上。网络和跨机房问题。如果消费者和Broker不在同一机房网络往返时间会直接体现在端到端延迟上。开发环境不明显生产环境要尽量避免跨地域消费。单机吞吐到了瓶颈。一个Broker的吞吐再高也有上限消息量大到一定程度单节点的分区Leader会成为热点。这时候就需要往集群方向扩展了——加Broker节点重新调整分区和副本分布。另外热词里有人搜“Kafka如何延迟30分钟消费”这里多说一句Kafka原生不支持按时间延迟消费实现这类需求一般有两种思路一种是消息里带时间戳消费者轮询时判断是否到期没到期就重新放回队列另一种是引入延迟队列组件比如RocketMQ原生支持延迟消息或者在Spring Boot里用定时任务扫描暂存表到期再投递到消费者。这个属于进阶设计系列后面我会单独写。5.4 可视化工具连接不上 / 集群模式注意点可视化工具连不上的问题90%和advertised.listeners有关。无论用Offset Explorer还是Kafka UI连接地址填了localhost:9092但Broker配置的advertised.listeners写的是别的地址工具就会连接失败。处理方式就一句话保证客户端访问地址和Broker宣告地址一致。如果做集群还有几个要注意的配置项每个Broker的broker.id必须唯一。集群模式下zookeeper.connect要指向同一个Zookeeper集群KRaft模式则是同一个Controller集群。所有Broker的advertised.listeners要能被其他Broker和客户端访问到否则集群会反复出现Leader选举失败的报错。Topic的副本因子不能大于Broker数量比如3个Broker的集群副本因子最多设3。我从一开始就被Kafka的“报错复杂”吓过后来发现大多数问题就是地址配置、端口不通、资源不足这三类抓住这三条能解决八成以上搭建问题。最后再分享一个我自己在工作中养成的习惯每次搭完Kafka我做的第一件事不是写业务代码而是先跑一遍生产者、消费者命令行脚本再用Offset Explorer确认Topic、Offset都正常显示。这个验证链路只要通了后面接入Spring Boot或者其他客户端时问题范围就能大幅缩小——至少你可以确信是业务代码的问题而不是Kafka服务本身的问题。搭建这一步做到这个确认程度就真的算牢固了。
返回列表