
如果你平时写的是HTTP接口忽然接到需求说要把订单消息丢给Kafka再由下游系统消费处理刚开始多半会有点发怵。Spring Boot集成Kafka这件事网上教程一搜一大把但很多要么只贴代码不讲原理要么版本太老跑不起来真正能把环境搭建、核心配置、代码演示和问题排查串成一条完整链路讲清楚的并不多。我最早做这块的时候也踩了不少坑版本对不上、消费者起不来、消息发出去却消费不到每个问题都能耗掉半天。这篇文章就基于我在实际项目里的使用经验从零开始带你把Spring Boot集成Kafka的完整流程跑通包含本地环境搭建、核心概念速记、生产者消费者代码实现、常见问题排查全程实测可复现建议先收藏再动手。1. 项目整体设计为什么选Kafka以及集成方案的取舍1.1 Kafka是消息队列但又不只是消息队列做后端这些年我见过太多团队在引入中间件时根本没想清楚“为什么”。Kafka的定位其实是分布式的消息流平台它不仅能当消息队列用还能做事件流处理、日志收集、数据管道。但日常业务里我们最常用的还是它的消息队列能力在订单系统、库存系统、积分系统之间架一道异步的缓冲层。我举个实际场景。传统单体应用里订单创建成功后要调积分接口、短信接口、搜索同步接口一个接口挂了整个下单链路就卡住。引入Kafka之后订单服务只需要把“订单已创建”这个事件写进Topic积分、短信、搜索各自异步去消费互不阻塞谁挂了也不影响别人。这个解耦效果用同步HTTP调用是实现不了的。既然说解耦就绕不开Kafka和RabbitMQ的对比。很多人刚开始选型会纠结我给个直观的参考维度KafkaRabbitMQ吞吐量极高适合海量日志、事件流中等适合业务消息路由消息模型以Topic维度消费者主动拉取以Queue维度支持多种路由模式消息回溯支持按offset重新消费天然实现消费后默认删除不太方便回溯可靠性多副本持久化配合ACK处理消息确认机制灵活典型场景大数据管道、异步解耦、日志采集传统业务通知、任务分发结论很简单如果你的核心诉求是“把多个下游系统解耦同时保证高吞吐”选Kafka基本不会错。1.2 Spring Boot集成Kafka的三种姿势Spring Boot项目接入Kafka主流有三条路。第一种是用Spring官方提供的spring-kafka框架。它把Kafka原生客户端包了一层提供KafkaListener注解、KafkaTemplate模板类、自动配置和错误处理机制什么版本兼容性、连接工厂、序列化器配置框架都给你打理好了。绝大多数项目用这个就够了这也是我推荐的默认选择。第二种是用Kafka原生客户端。引入kafka-clients依赖自己new KafkaProducer和KafkaConsumer。好处是没有任何封装魔法每一步都看得见摸得着。坏处也很明显所有模板代码都要自己写连接管理、线程调度、提交offset都得自己把控开发效率和可维护性都差一些。除非你对Kafka底层有特殊控制需求否则不推荐在Spring Boot项目里这么干。第三种是Spring Cloud Stream。它把Kafka、RabbitMQ这些消息中间件抽象成统一的编程模型用Binder概念屏蔽底层差异。听起来很美但实际上抽象层越多排查问题越费劲配置也绕。我见过不少团队引入Spring Cloud Stream后光调试消息转换格式就花了一周。除非你有“一套代码同时对接多种消息中间件”的强需求否则我建议绕开它。所以这次的demo就用spring-kafka原因很简单它沉淀了Spring社区对消息场景的最佳实践配置直观、代码量少出问题也容易在社区找到答案。1.3 本次demo的架构与代码结构我们先明确一下这个演示项目要做什么一个普通的Spring Boot应用既扮演生产者通过HTTP接口往Kafka发送消息也扮演消费者监听指定Topic并处理消息。实际生产里生产者和消费者大概率是两个不同服务但学习阶段放在同一个工程里最直观你能一眼看到消息从发送到消费的完整链路。工程结构我按下面这样规划config用于存放Kafka相关配置类比如把序列化和反序列化配置补充完整controller提供HTTP接口触发消息发送service封装KafkaTemplate的发送逻辑listener用KafkaListener编写消费者监听逻辑整体架构就两部分Kafka服务端用单机KRaft模式Spring Boot应用通过spring-kafka提供的Template和Listener完成收发。这个组合复现起来最简单也最贴近日常开发的主干路径。2. 环境搭建本地Kafka从下载到跑通2.1 四个核心概念先过一遍动手安装之前先把几个概念打通否则后面配置你只会照抄但不懂含义。Topic是消息的分类你可以理解为数据库里的表。生产者和消费者都是围绕Topic进行读写。Partition是Topic的物理分片一个Topic可以拆成多个Partition每个Partition内部消息有序。Partition数量直接决定了消费并行度一个Partition在同一时刻只能被一个消费者实例消费。Offset是消费者在分区内的读取位置。消费者每消费一条消息offset就往后移一位。Kafka不删除已消费的消息而是通过offset记录消费到哪里这也是它能实现消息回溯的根本原因。Consumer Group是逻辑上的消费者集合。同一个Group里的多个消费者分摊消费一个Topic下的所有Partition一条消息只会被同一个Group里的一个消费者实例处理。但不同Group之间各自独立每个Group都能把消息完整消费一遍。你只需要记住一句话Topic是消息的家Partition是家里的书柜Offset是书签Consumer Group是看书的人。理解了这个模型后面配置和排查问题就顺了。2.2 下载与版本选择Kafka的下载很简单去Apache官网下载编译好的二进制包就行。我这次用的是Kafka 3.6.0下载包名大概是kafka_2.13-3.6.0.tgz。这里有个容易搞混的点包名里的2.13是指编译Kafka所用的Scala版本Kafka自身是Java写的这个Scala版本只跟Kafka的构建工具链有关跟我们选型没关系。真正要关心的是Kafka的版本号和你在Spring Boot里用的Spring Kafka版本是否兼容。版本选择我给个实用建议Spring Boot 2.7项目自带Spring Kafka 2.8建议配Kafka 3.x都能跑Spring Boot 3.x要求JDK17Spring Kafka是3.x系列API上有一些调整但基础用法区别不大。如果你还是JDK8的老项目就别强行上Spring Boot 3老老实实用2.7加上Kafka 3.x完全够用。提醒一句下载前确认本机Java环境Kafka服务端要求JDK8或以上实测JDK8和JDK17都没问题。2.3 用KRaft模式单机启动Kafka以前启动Kafka必须先启动ZooKeeper命令长、配置多非常麻烦。从Kafka 2.8开始引入了KRaft模式把元数据管理直接做进了Kafka内部Kafka 3.x给了ZooKeeper模式最后的时间到了4.x就完全移除依赖了。所以我们直接用KRaft模式搭建单机环境步骤少很多。我以Linux/Mac命令为例Windows把.sh后缀换成.bat就行。第一步解压并进入目录tar -xzf kafka_2.13-3.6.0.tgz cd kafka_2.13-3.6.0第二步生成集群唯一ID。KRaft用这个ID初始化存储目录kafka-storage.sh random-uuid执行后会输出一串UUID比如vGx3mX0WQTKQqXmZmZlQfQ。先复制下来待用。第三步格式化存储目录kafka-storage.sh format -t 上面生成的UUID -c config/kraft/server.properties格式化时终端会提示This tool will reformat the log directory确认输入y回车即可。这个步骤相当于给Kafka初始化了一个“数据盘”不执行的话启动会直接报错。第四步启动Kafka服务kafka-server-start.sh config/kraft/server.properties看到日志里出现[/tmp/kraft-combined-logs] is now available和Kafka Server started字样就说明服务起来了。整个流程不到一分钟远没有以前搭ZooKeeper那么痛苦。2.4 用命令行验证Kafka基本可用性服务端跑起来以后先别急着写代码用命令行把链路验证一遍。这样后面Spring Boot联调如果失败你至少能确定问题不在Kafka环境本身。创建一个演示用的Topic分区设3个副本因子设1。单机环境副本因子只能设1设多了会报错kafka-topics.sh --create --topic demo-topic --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1查看Topic描述能确认3个分区都正常kafka-topics.sh --describe --topic demo-topic --bootstrap-server localhost:9092然后开一个终端运行控制台消费者kafka-console-consumer.sh --topic demo-topic --from-beginning --bootstrap-server localhost:9092再开另一个终端运行控制台生产者kafka-console-producer.sh --topic demo-topic --bootstrap-server localhost:9092在生产者终端输入hello kafka切到消费者终端能看到同样内容打印出来说明Kafka服务端完全正常。这一步跑通了我们就可以专心写Spring Boot代码了。3. Spring Boot集成Kafka代码实操与完整演示3.1 新建工程与依赖引入我先用Spring Boot 2.7.18版本建了一个空工程这也是当前JDK8项目里很常见的版本组合。如果你用的是JDK17直接上Spring Boot 3.x也完全没问题核心代码差异不大。pom.xml里至少要引入两个依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId /dependencystarter-web是为了暴露HTTP接口做发送入口spring-kafka则是本次的核心依赖。版本号不用手动填Spring Boot的BOM会帮我们管理好对应的Spring Kafka版本。这里我多说一句很多人喜欢在pom里强行指定spring-kafka版本如果和Spring Boot不匹配就容易出现一堆莫名其妙的NoSuchMethodError所以默认交给BOM管理最稳妥。3.2 理解并配置application.ymlSpring Boot集成Kafka的配置全都可以写在application.yml里核心部分我拆开讲。spring: kafka: bootstrap-servers: localhost:9092 producer: retries: 3 acks: all key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer consumer: group-id: demo-group auto-offset-reset: earliest enable-auto-commit: false key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer listener: ack-mode: manual_immediatebootstrap-servers是Kafka服务端地址多个节点用逗号分隔。producer部分的retries表示发送失败后的重试次数acks配置all表示消息要等到所有副本都写入成功才算完成这是保证消息不丢的关键配置。consumer部分需要注意两个参数。auto-offset-reset取值有earliest、latest、none三种earliest表示消费者组没有提交过offset时从最早的消息开始消费latest表示只消费启动之后新到的消息none表示没有offset直接报错。第一次调试建议用earliest否则你发完消息再启动消费者可能一条都消费不到。enable-auto-commit我配置成false关闭自动提交。这是生产环境推荐的姿势配合listener的ack-mode手动确认能避免很多丢消息和重复消费的问题。这些都是我踩过坑之后的总结后面排查环节还会展开讲。3.3 生产者代码实现生产者核心是KafkaTemplate它是Spring Kafka最常用的消息发送模板线程安全可以在Service里注入一次反复使用。Component public class KafkaProducerService { private static final Logger log LoggerFactory.getLogger(KafkaProducerService.class); Resource private KafkaTemplateString, String kafkaTemplate; public void send(String topic, String message) { kafkaTemplate.send(topic, message).whenComplete((result, ex) - { if (ex ! null) { log.error(消息发送失败, topic{}, message{}, topic, message, ex); } else { RecordMetadata metadata result.getRecordMetadata(); log.info(消息发送成功, topic{}, partition{}, offset{}, topic, metadata.partition(), metadata.offset()); } }); } }这里我用whenComplete异步回调处理发送结果而不只是调send之后就不管了。KafkaTemplate的send本身是异步的立即返回的是Future。如果你不关心结果线上消息发失败了你都不知道。记录分区和offset信息对于排查问题非常有用比如你怀疑某个key的消息没有进到预期的分区看日志就清楚了。发送时还可以指定key比如kafkaTemplate.send(topic, orderId-12345, message)。Kafka会默认对key做hash相同key会路由到同一个分区。如果你希望某个业务ID的消息严格有序就一定要带key去发送。3.4 消费者实现与手动ACK消费者用KafkaListener注解挂在方法上框架会自动监听指定Topic并把消息传入方法。我们这里配置了手动ACK模式所以方法参数里加上Acknowledgment。Component public class KafkaConsumerListener { private static final Logger log LoggerFactory.getLogger(KafkaConsumerListener.class); KafkaListener(topics demo-topic, groupId demo-group) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { try { log.info(收到消息, partition{}, offset{}, value{}, record.partition(), record.offset(), record.value()); // 在这里处理业务比如写库里、调接口、更新缓存 ack.acknowledge(); } catch (Exception e) { log.error(消息处理失败, offset{}, record.offset(), e); // 按实际业务决定抛出异常还是记录到死信队列这里先记录日志 } } }手动ACK意味着offset的提交完全由你控制。a里ack.acknowledge()调用后Kafka才会把当前消费位置标记为已处理。这样处理逻辑真正成功之后再告诉Kafka“这条我搞定了”如果处理过程中抛异常offset未被提交下次启动会重新消费这条消息。当然这也引出了重复消费的问题后面排查章节会专门讲幂等。有些场景下业务处理耗时长担心消费者被Kafka判定为失联这个时候可以配合配置max.poll.interval.ms。但这个参数默认5分钟一般业务处理不会达到这个阈值真到了说明逻辑有问题先考虑异步化而不是无脑调参数。3.5 对外接口与联调演示为了演示方便我写一个简单的Controller用HTTP POST触发消息发送RestController RequestMapping(/kafka) public class KafkaController { Resource private KafkaProducerService producerService; PostMapping(/send) public String send(RequestParam String topic, RequestParam String message) { producerService.send(topic, message); return send success; } }然后启动Spring Boot应用用curl做一次发送curl -X POST http://localhost:8080/kafka/send?topicdemo-topicmessagehello-springboot-kafka观察控制台日志你会先看到生产者回调日志“消息发送成功partition1, offset15”紧接着消费者监听日志就出来了“收到消息partition1, offset15, valuehello-springboot-kafka”。链路完全跑通。这里我建议你也试一下用命令行生产者给Spring Boot消费者发消息或者反过来用Spring Boot生产者给命令行消费者发消息交叉验证两边都没问题。我之前联调时就遇到过Spring Boot能给自己发、能给自己收但接不上外部系统的情况结果是对面消费者用的groupId和我们不一致互相都收不到。3.6 进阶事务消息、批量消费、JSON消息基础链路跑通之后有几个高频的进阶需求值得提前了解。事务消息。如果你希望消息发送和本地数据库事务保持一致性可以给KafkaTemplate开事务。先配置spring.kafka.producer.transaction-id-prefix再往启动类或配置类加EnableKafka然后发送方法上标TransactionalTransactional public void sendWithTransaction(String topic, String message) { kafkaTemplate.send(topic, message); }这种模式适合“本地库操作和发消息要么都成功要么都失败”的场景。但它对服务端也有要求事务ID前缀必须全局唯一也会带来一定性能损耗没有强一致需求的话不必轻易上。批量消费。默认一条一条处理可能吞吐不够可以把监听方法改成批量模式。在yaml里加上spring.kafka.listener.typebatch方法入参改成ListKafkaListener(topics demo-topic, groupId demo-group) public void onBatchMessage(ListConsumerRecordString, String records) { log.info(批量收到{}条消息, records.size()); // 批量处理逻辑 }JSON消息。生产者和消费者传String虽然简单通用但业务上更常见的是传对象。一种做法还是传String把对象序列化成JSON字符串消费者用ObjectMapper解析这样序列化器不绑定灵活可控。另一种做法是配置JsonSerializer和JsonDeserializer让框架直接帮我们转对象。我建议核心场景用String加自己解析跨系统联调时最不容易翻车。框架的JSON转换器看起来省事但不同服务的Jackson版本差异很容易搞得你头大。4. 常见问题与排查技巧实录4.1 连接失败先分清楚客户端还是服务端这是新手遇到最多的问题报错通常是Connection refused或者Timed out。排查思路很简单先确认Kafka服务端到底有没有起来。本机执行jps看有没有Kafka进程。没有就检查启动日志有没有报错。确定进程存在后再用telnet测试端口通不通telnet localhost 9092如果telnet不通多半是监听地址配置问题。默认server.properties里listeners是PLAINTEXT://:9092正常可以用localhost访问。如果你改过配置或者Kafka跑在Docker里要额外检查端口映射和主机名配置。还有一个容易忽略的点客户端配置的bootstrap-servers值和Kafka广告地址不一致。如果Kafka集群做了内外网分离broker把advertised.listeners设成内网地址而客户端从公网访问就会出现“端口是通的但连接不到broker”的诡异现象。本地单机环境一般遇不到但多节点或容器环境要特别注意。4.2 消息发出去却消费不到消息发送日志显示成功但消费者日志一条都没有这个问题我遇到太多次了。绝大多数原因有三个。第一groupId不一致。消费者监听方法里如果写了groupId会覆盖yaml里的配置。你生产者在demo-group消费者却是another-group两边各自独立消费者当然只能消费到启动之后新来的消息。检查方式就是把两边的groupId对齐。第二auto-offset-reset配成latest消费者在消息发送之后才启动那么之前发的消息就全部被跳过了。第一次联调建议改成earliest或者干脆用一个新的groupId测试。第三消费者组的历史offset已经提交过。同一个groupId之前消费过这个TopicKafka会把offset记在内部里下次启动会从上次的位置继续新消息来之前什么都收不到。想强制从头消费可以给消费者换一个全新groupId。排查时用命令行看消费组的offset情况很管用kafka-consumer-groups.sh --describe --group demo-group --bootstrap-server localhost:9092这个命令能显示每个分区的当前偏移量、LOG-END-OFFSET和消费者Lag值。如果Lag长期不为0但消费者端没日志说明消息在被消费但可能抛异常被吞了优先查消费逻辑里有没有catch后静默处理。4.3 反序列化报错序列化器不匹配报错一般长这样org.apache.kafka.common.errors.SerializationException: Cant deserialize data。原因十有八九是生产者的value-serializer配的是StringSerializer发送了字节数据或JSON字节但消费者把value-deserializer配成了JsonDeserializer两边不匹配。排查思路是先把配置对齐。如果你消息体是一个JSON字符串客户端两边都写成String系列化器然后在代码里自己解析JSON这是最简单粗暴又最不容易出错的做法。一旦你用JsonSerializer框架默认会往消息头里写入类型信息而接收方的包名类名跟发送方不完全一致就会抛类型转换异常。所以我的建议是项目内部可以统一JSON序列化器但跨系统联调最稳妥还是String加显式解析。4.4 重复消费与消息丢失手动ACK是一把双刃剑开启手动ACK之后重复消费和消息丢失看起来都被避免了但处理不当会引入新问题。先说重复消费。业务处理成功了但是ack.acknowledge()还没来得及提交应用就宕机了offset没更新重启后这条消息还会被再次消费。或者处理完业务返回给调用方失败了你catch住异常没有抛出去但也没ack下次轮询还是会再拿一次。这个时候消费者侧要保持幂等每条消息都带一个业务唯一ID消费时先查这个ID有没有处理过处理过就直接跳过。哪怕是简单的把唯一ID存到数据库加唯一索引都能挡住大量重复消息。再说消息丢失。如果把ack写在了业务处理之前或者开启自动提交而处理超时都会出现“消息已经标记消费但业务没完成”的情况。手动ACK的正确姿势是业务处理成功后再ack处理失败要明确记录并安排补偿千万不要把ack放在try块第一行这是我见过很多同学踩过的坑。4.5 延迟高、吞吐上不去先看分区数再看并发消息消费延迟高很多人的第一反应是调消费者并发。但有个前提分区数决定了消费者并发上限。一个Partition同一时刻只能被同一个Group里的一个消费者实例消费如果你的Topic只有一个Partition消费者开20个线程也只有1个在真正消费其余全部闲置。所以生产环境创建Topic时就要评估好分区数。分区多一方面提高并发另一方面也带来分区副本选举和消息顺序的问题不能一味贪多。一般经验是分区数大于等于消费组内消费者实例数的整数倍保证每个消费者都有分区可拉并且尽量均衡。消费端还可以调几个参数max.poll.records控制单次拉取消息条数拉太多处理太慢会触发rebalancemax.poll.interval.ms控制两次poll的最大间隔超过会被认为消费者失联concurrency配置并发消费线程数。这些参数配合partition数量合理调优延迟通常能降下来。我之前遇到过一个延迟积压严重的情况就是把max.poll.records从500降到50处理时间立刻缩短rebalance也少了。4.6 Spring Boot与Spring Kafka版本匹配问题最后说一个隐蔽的坑。Kafka客户端和服务端有版本兼容问题Kafka 3.x配合旧版spring-kafka可能报UnsupportedVersionException这是因为客户端版本低于broker的协议版本。解决思路很简单Spring Boot版本会管理对应Spring Kafka版本除非有特殊需求不要手动改版本。比如Spring Boot 2.7.x自带Spring Kafka 2.8.x它支持的broker版本范围本来就包含Kafka 3.x实测没问题。如果你引入Spring Kafka 3.x但Spring Boot还是2.x可能会因为Spring框架版本过低导致启动报错。另外注意Spring Boot 3.x要求JDK17如果你的团队还在用JDK8别为了追新硬上Spring Boot 3。技术选型永远是稳定优先Kafka版本只是工具能稳定跑通业务才是目的。最后再分享一个我个人的实操习惯每次搭Kafka相关项目我都会先在命令行把生产者和消费者手动跑通再动Spring Boot代码。这样做一次能帮你省掉后面至少三个小时的排查时间因为环境问题已经被提前排除了。另外本地开发强烈建议用KRaft模式单机启动不要再走ZooKeeper的老路配置简单不说换机器重新搭建也快。等你把这条链路跑通了再逐步加上事务消息、批量消费、监控告警业务接进来自然就顺了。