ARTICLE DETAIL

资讯详情

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

SpringBoot集成Kafka实战:从环境搭建到消息收发全流程解析

SpringBoot集成Kafka实战:从环境搭建到消息收发全流程解析 最近好几拨朋友都在问SpringBoot集成Kafka怎么做网上的教程要么停留在ZooKeeper时代要么贴一堆配置不给解释照着做还经常跑不通。这篇文章我从环境搭建开始一步步演示SpringBoot怎么集成Kafka把配置项、核心代码、展示效果、常见坑全部理清楚保证你跟着走完就能在自己项目里用起来。这套方案适合正在用SpringBoot做微服务、需要引入消息队列做异步解耦的Java后端同学也适合准备自己搭一套Kafka环境深入学习原理的人。整个流程我按生产环境可参考的标准来组织但演示部分又足够简单哪怕你只是第一次接触Kafka也不用担心跟不上。1. 为什么要用Kafka以及整套方案的设计思路1.1 Kafka在技术栈里的位置消息队列在系统里扮演的核心角色就三个解耦、削峰、异步。Kafka的特点在于吞吐量极高、数据可持久化、天然支持分布式和水平扩展因此在日志采集、用户行为追踪、订单状态流转、系统间数据同步这些场景下特别稳。相比RabbitMQKafka更适合大数据量的流式处理而SpringBoot作为Java微服务的标配框架和Kafka客户端集成非常顺滑依赖一加、配置一写生产者和消费者就能跑起来。我见过很多团队在初期直接用Feign调接口或者自己写线程池处理异步任务一旦流量上来线程池队列堆满、接口超时、任务丢失这些问题全冒出来了。换成Kafka之后生产者只管发送消息消费者按自己的能力拉取处理两边互不拖累这也是我在项目里最常用Kafka的原因。1.2 版本选择与架构演进Kafka目前已经全面进入KRaft模式也就是不再依赖ZooKeeper来管理元数据了。早期版本架构里Kafka Broker、Topic、分区等元数据都存放在ZooKeeper上导致运维需要同时维护两套集群。从Kafka 3.3开始KRaft模式进入生产可用状态Kafka 3.5之后主流版本基本都推荐直接用KRaft模式部署单节点轻量测试尤其省事。SpringBoot对Kafka的支持主要在spring-kafka这个项目中SpringBoot 2.x默认管理的是Kafka客户端2.8.x左右SpringBoot 3.x会管理更新的客户端版本。我建议你新项目直接上SpringBoot 3.x加Kafka 3.6以上版本毕竟新项目没必要在旧版本上纠结而且新版客户端对性能协议都有优化。提示如果你还在用SpringBoot 2.xKafka版本选3.x客户端也是兼容的只要注意Broker和客户端之间的小版本差异不要太大即可一般2.8到3.6之间兼容性都还行。1.3 这套演示项目的目标与功能范围我这套演示项目要做的功能很简单通过一个HTTP接口发送用户注册消息到Kafka消费者监听到消息之后模拟保存日志和处理业务。麻雀虽小五脏俱全里面包含了工程搭建、Topic管理、生产者发送、消费者监听、序列化配置、参数调优这些完整链路。你照着搭建完之后把业务逻辑替换成自己的场景就行比如下单通知、短信发送、埋点上报思路完全一致。2. 环境准备先把Kafka跑起来2.1 下载Kafka与基础环境要求需要先装好JDK建议JDK 8以上我本机用的是JDK 17。Kafka的官方压缩包自带启动脚本不需要额外安装什么数据库或者中间件真实执行下载的是编译好的二进制包。去Kafka官网下载页选择最新的二进制版本比如kafka_2.13-3.6.2.tgz其中2.13是配套的Scala编译器版本3.6.2才是Kafka真实版本。有人会纠结选哪个Scala版本其实对普通使用没影响选官方推荐的就行。下载完成后解压到一个没有空格的路径比如/opt/kafka或者D:\kafka解压完目录结构大致如下kafka_2.13-3.6.2/ ├── bin/ # 启动脚本 ├── config/ # 配置文件 ├── libs/ # 依赖包 └── logs/ # 日志目录2.2 使用KRaft模式单机启动Kafka进入Kafka目录先修改config/kraft/server.properties里这几个关键参数process.rolesbroker,controller node.id1 controller.quorum.voters1localhost:9093 listenersPLAINTEXT://:9092,CONTROLLER://:9093 advertised.listenersPLAINTEXT://localhost:9092 log.dirs/tmp/kraft-combined-logs对照说明一下process.roles表示这个节点同时充当broker和controller角色单机模式就这么写。controller.quorum.voters配置的是集群中controller节点的地址格式是node.idhost:port这里固定端口9093。listeners是服务监听地址9092给客户端连接用9093给controller通信用。advertised.listeners是告诉客户端连接哪个地址这个特别关键如果你要让其他机器访问必须改成对应的IP或域名。log.dirs是数据存储目录生产环境一定要改成独立的数据盘路径不要用/tmp。第一次启动KRaft模式需要生成集群ID并格式化存储目录直接执行官方提供的一条命令KAFKA_CLUSTER_ID$(bin/kafka-storage.sh random-uuid) bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.propertiesWindows环境要用bin\windows\下的bat脚本命令逻辑一样。格式化完成之后启动Kafkabin/kafka-server-start.sh config/kraft/server.properties看到类似“Kafka Server started”的日志就说明启动成功了。这个单机模式不需要再启动ZooKeeper省掉一个进程比老教程简单很多。2.3 创建Topic并验证消息收发Kafka里面的Topic相当于一个消息分类生产者往Topic里写消息消费者从Topic里读消息。我演示用的Topic叫user-topic采用3个分区。分区数决定了消息的并行处理能力生产环境一般结合消费者数量和吞吐量来设计这里演示设3个就够了。bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic user-topic --partitions 3 --replication-factor 1查看Topic列表验证是否创建成功bin/kafka-topics.sh --bootstrap-server localhost:9092 --list启动自带的生产者控制台手动发送一条消息测试一下bin/kafka-console-producer.sh --broker-list localhost:9092 --topic user-topic输入一条文本然后回车再启动消费者控制台查看消息bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic user-topic --from-beginning能收到刚才输入的内容说明Kafka本身已经通了接下来就交给SpringBoot集成。2.4 初始化SpringBoot项目我用的开发工具是IntelliJ IDEA创建一个Spring Initializr项目Java版本选17依赖这里只需要加上Spring Web和Spring for Apache Kafka。如果你习惯用Maven核心依赖在pom.xml里就这两块dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId /dependencySpringBoot的自动配置会读取application.yml里spring.kafka开头的配置自动创建KafkaTemplate和ConsumerFactory这些对象省去很多手写代码。3. 核心代码实现生产者和消费者这样写3.1 配置文件里的关键参数先看application.yml我按生产可扩展的风格来写spring: application: name: kafka-demo kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.springframework.kafka.support.serializer.JsonSerializer acks: all retries: 3 properties: enable.idempotence: true linger.ms: 5 batch.size: 16384 consumer: group-id: user-group auto-offset-reset: earliest key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer properties: spring.json.trusted.packages: * spring.json.value.default.type: com.example.kafkademo.dto.UserMessage max.poll.records: 50 listener: concurrency: 3 missing-topics-fatal: false逐个说下关键参数的意义bootstrap-servers就是Kafka的地址列表生产环境集群就写多个节点用逗号隔开比如192.168.1.10:9092,192.168.1.11:9092。producer里key-serializer和value-serializer是指消息的key和value如何序列化。key直接用StringSerializer就够了value这里我用JsonSerializer意味着我们发送对象时会自动转成JSON。acks: all表示所有副本都确认后才算发送成功这是最强的可靠性保证虽然会损失一点延迟但数据安全性高。retries重试次数设为3配合enable.idempotence: true可以避免网络抖动导致的重复消息。linger.ms是生产者发送消息前等待更多消息进入批次的时间设5毫秒意味着攒一批再发能显著提升吞吐。batch.size是批次大小16KB适合大多数业务场景。consumer里group-id是消费者组ID同一个组内的消费者分摊消息。auto-offset-reset: earliest表示消费者没有提交offset时从头开始消费如果你的业务只关心新消息可以改成latest。这里最关键的坑在于JsonDeserializer反序列化。消费者拿到的是JSON字符串要转成什么类型JVM必须知道所以要么在配置里指定spring.json.value.default.type要么在注解里指定目标类型。我用配置方式把这个类写好类型信息就能正确还原。listener.concurrency是监听器的并发线程数这里设3配合Topic的3个分区可以实现3个线程并行消费。注意这个并发数不要大于分区数否则多余线程会空闲。3.2 定义消息实体类既然是用户注册的消息我定义一个UserMessage类package com.example.kafkademo.dto; import java.time.LocalDateTime; public class UserMessage { private Long userId; private String username; private String email; private LocalDateTime registerTime; public UserMessage() { } public UserMessage(Long userId, String username, String email, LocalDateTime registerTime) { this.userId userId; this.username username; this.email email; this.registerTime registerTime; } public Long getUserId() { return userId; } public void setUserId(Long userId) { this.userId userId; } public String getUsername() { return username; } public void setUsername(String username) { this.username username; } public String getEmail() { return email; } public void setEmail(String email) { this.email email; } public LocalDateTime getRegisterTime() { return registerTime; } public void setRegisterTime(LocalDateTime registerTime) { this.registerTime registerTime; } Override public String toString() { return UserMessage{ userId userId , username username \ , email email \ , registerTime registerTime }; } }注意一定要保留无参构造函数。JSON反序列化时需要先创建对象再填充字段没有无参构造会直接报错。我在实际开发中经常看到有人写了带参构造就忘掉无参构造结果运行期才暴露问题排查起来很心累。3.3 生产者Producer与发送接口SpringBoot的KafkaTemplate封装了所有发送逻辑我们只需要注入后调用send方法它默认会使用配置里的序列化器和acks策略。我写一个发送服务package com.example.kafkademo.service; import com.example.kafkademo.dto.UserMessage; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.SendResult; import org.springframework.stereotype.Service; import java.util.concurrent.CompletableFuture; Service public class KafkaProducerService { private final KafkaTemplateString, Object kafkaTemplate; public KafkaProducerService(KafkaTemplateString, Object kafkaTemplate) { this.kafkaTemplate kafkaTemplate; } public boolean sendUserMessage(UserMessage message) { String key String.valueOf(message.getUserId()); CompletableFutureSendResultString, Object future kafkaTemplate.send(user-topic, key, message); future.whenComplete((result, ex) - { if (ex null) { System.out.println(消息发送成功: message , offset: result.getRecordMetadata().offset()); } else { System.err.println(消息发送失败: ex.getMessage()); } }); return true; } }这里有几个细节值得注意。第一是send方法传入key相同key的消息会进同一个分区从而保证同一用户的注册消息严格有序。第二是send方法返回CompletableFuture异步回调里能拿到发送结果和offset生产环境强烈建议像这样处理发送结果而不是发完就不管了否则消息丢失很难定位。第三是我把KafkaTemplate的泛型定义成String和Object这样同一套模板可以发送不同类型的对象灵活性更高。再写一个Controller暴露HTTP接口package com.example.kafkademo.controller; import com.example.kafkademo.dto.UserMessage; import com.example.kafkademo.service.KafkaProducerService; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import java.time.LocalDateTime; RestController RequestMapping(/api/user) public class UserController { private final KafkaProducerService producerService; public UserController(KafkaProducerService producerService) { this.producerService producerService; } GetMapping(/register) public String register(RequestParam Long userId, RequestParam String username) { UserMessage message new UserMessage(userId, username, username example.com, LocalDateTime.now()); producerService.sendUserMessage(message); return 消息已发送; } }我故意用GET接口方便浏览器直接测试实际项目中你肯定要换成POST并加参数校验这里重点看Kafka集成接口设计从简。3.4 消费者Consumer与监听逻辑消费者是SpringKafka里最爽的部分只用KafkaListener注解就可以监听一个或多个Topicpackage com.example.kafkademo.consumer; import com.example.kafkademo.dto.UserMessage; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; Component public class UserMessageConsumer { KafkaListener(topics user-topic, groupId user-group) public void onMessage(UserMessage message) { System.out.println(收到用户注册消息: message); // 这里可以写你的业务逻辑比如保存日志、发送邮件、更新统计等 } }注意这里方法参数直接写UserMessage类型spring-kafka会自动完成反序列化但前提是我们在配置文件里指定了spring.json.value.default.type。如果消息发送时用的不是JSON序列化器或者消息体格式不匹配这里就会报反序列化异常。3.5 配置监听器容器工厂进阶调整默认情况下KafkaListener使用SpringBoot自动配置的容器工厂大部分场景不需要改。但如果要做批量消费监听、错误处理器定制、并发调整就需要自己定义了一个监听容器的工厂。我建议你在demo阶段先不改跑通了再按需调整。真的需要批量处理时可以在启动类加EnableKafka注解然后定义一个ConcurrentKafkaListenerContainerFactoryBean public ConcurrentKafkaListenerContainerFactoryString, Object kafkaListenerContainerFactory( ConsumerFactoryString, Object consumerFactory) { ConcurrentKafkaListenerContainerFactoryString, Object factory new ConcurrentKafkaListenerContainerFactory(); factory.setConsumerFactory(consumerFactory); factory.setConcurrency(3); factory.setBatchListener(true); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); return factory; }setBatchListener(true)开启批量模式后监听方法参数需要改成List 一次拉取一批消息。MANUAL_IMMEDIATE配合手动提交offset等业务处理成功后再提交防止处理中途崩溃导致offset丢失。如果业务要求“消息必须不丢”手动提交是必选项默认的自动提交在某些异常场景下会丢消息。4. 完整演示从发送到消费的闭环跑通4.1 启动项目并观察日志先确认Kafka服务在运行然后在IDEA里启动SpringBoot应用。启动日志里如果能看到类似“Started KafkaDemoApplication”和Spring Kafka相关的初始化信息说明自动配置生效了。4.2 通过HTTP接口发送消息浏览器地址栏输入http://localhost:8080/api/user/register?userId1001usernamezhangsan页面返回“消息已发送”之后立刻回看IDEA控制台你会看到类似这样的日志收到用户注册消息: UserMessage{userId1001, usernamezhangsan, emailzhangsanexample.com, registerTime2025-01-15T10:23:45} 消息发送成功: UserMessage{...}, offset: 0发送成功回调和消费日志不一定谁先出现因为发送是异步的消费者收到消息后立刻打印几乎同时发生。能同时看到这两行说明Producer到Broker再到Consumer的链路已经完整打穿。4.3 验证消息在Kafka里的落盘情况再打开一个消费者终端bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic user-topic --from-beginning如果你配置的是JsonSerializer这里控制台会显示一段JSON而不是Java对象的toString格式。这正好验证了消息在Kafka中存储的介质是JSON文本消费者反序列化时才还原成对象。4.4 编写JUnit测试类代替HTTP接口有时候不想启动Web服务只想快速验证Kafka链路可以写个测试类package com.example.kafkademo; import com.example.kafkademo.dto.UserMessage; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.kafka.core.KafkaTemplate; import java.time.LocalDateTime; SpringBootTest class KafkaDemoApplicationTests { Autowired private KafkaTemplateString, Object kafkaTemplate; Test void sendMessageToKafka() { UserMessage message new UserMessage(2001L, lisi, lisiexample.com, LocalDateTime.now()); kafkaTemplate.send(user-topic, String.valueOf(message.getUserId()), message); // 休眠几秒等待消费者处理 try { Thread.sleep(3000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }加不加SpringBootTest取决于你细到什么粒度。用测试类验证的好处是不依赖接口后续你写单元测试或者集成测试都能复用它。4.5 模拟消息积压和并发消费为了看到多分区并行消费的效果你可以在短时间内连发十几条消息然后观察控制台打印线程名称。如果并发配置生效会有多个不同的线程同时输出消费日志这就是Kafka水平扩展能力的直观体现。这条演示很容易受auto-offset-reset影响。如果你是先启动消费者、再发送消息无论配置earliest还是latest都收得到但如果你先把消息发到了Topic再启动消费者并且配置是latest消费者会从最新offset开始消费之前积压的消息全部收不到。很多新手在这里踩坑以为消息丢了其实是被offset策略跳过了。5. 常见问题与排查技巧实录5.1 Kafka连接失败启动项目报Connection refused出现频率最高的原因有两个。第一个是Kafka服务没有启动或者不是在本机的9092端口启动。第二个是配置的advertised.listeners不对尤其是当应用和Kafka不在同一台机器时如果advertised.listeners写的localhost客户端连接时依然会去连localhost自然连不上。排查方法很简单先用命令行工具验证bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092如果这个命令能正常拉取Broker版本信息说明Kafka服务可用问题去应用配置里找。注意要保证应用所在机器能访问Kafka主机的9092端口云服务器还需要在安全组里放行。5.2 消费者收不到消息消费者启动后控制台没有任何输出优先检查这些项检查消费者组的offset提交情况。使用命令查看bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group user-group输出里会有CURRENT-OFFSET和LOG-END-OFFSET两列如果CURRENT-OFFSET等于LOG-END-OFFSET说明消息已经被消费过了不会再触发监听。检查auto-offset-reset配置。新消费组首次启动时才会触发这个策略如果组已经存在且有已提交的offset它不会重新消费旧消息。检查Topic名称是否一致。消费者监听的是topicA发送者发到topicB自然收不到。这种低级错误我真见过不少两边字符串只要差一个字符就完蛋。检查监听器并发和分区数量关系。理论上分区数为3并发为3是有效的如果你把监听并发调成了10超出部分线程只能闲置并不会增加消费速度。5.3 反序列化报错最常见的报错是JsonMappingException或者TypeIdResolver相关异常。原因是消费者不知道JSON要反序列化成什么类型。我之前配了spring.json.value.default.type这里要说清楚这个配置只对JsonDeserializer默认生效。如果你在方法参数上指定类型Spring会优先根据参数类型来反序列化但依然需要配置不被信任的包。SpringKafka为了安全JsonDeserializer默认只信任java.util和java.lang包我们自己定义的com.example.kafkademo.dto就必须显式加入白名单。我配置文件里写了spring.json.trusted.packages: *这样最省事缺点是安全性弱一点。生产环境建议改成具体的包名比如com.example.kafkademo.dto。如果你压根不想用JSON也可以用StringSerializer和StringDeserializer。消息内容为JSON字符串消费者收到后用ObjectMapper解析。这种方式规避了类型绑定问题在很多公司里反而最常用因为消费逻辑可以自己对字符串做解析灵活度高。5.4 消息重复与幂等设计Kafka只能保证不丢消息不保证不重复消息。生产端的重试、消费端在重启后重新提交都可能导致一条消息被消费多次。我的做法是消费逻辑必须设计成幂等。最简单的幂等方案是在业务表里加一个唯一键直接用消息里的业务ID作为唯一键插入插入冲突就跳过。也可以在消费者里用一个Redis记录已处理的消息ID处理前检查是否出现过。5.5 常用排查命令速查表总结一下在线排查时最实用的Kafka命令场景命令查看Topic是否存在bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic user-topic从头重新消费消息创建带新group-id的消费者并设置earliest或删除消费组后重新监听查看消费组消费进度bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group user-group查看Broker是否可用bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092模拟生产者手动发消息bin/kafka-console-producer.sh --broker-list localhost:9092 --topic user-topic注意删除消费组的命令是kafka-consumer-groups.sh --delete --group user-group但前提是该组当前没有活跃消费者。演示环境随便用生产环境慎用换新group-id观察日志更安全。5.6 一些提高生产力的配置心得看完上面这些老读者应该已经发现SpringBoot集成Kafka最核心的技术点无非三条序列化方式、消费组策略、发送可靠性参数。我在多个项目里总结出一套比较稳的默认配置组合你直接抄作业也没问题消息对象统一用DTO不带业务逻辑序列化用JSON。生产端acks设all开启幂等配合retries防止瞬时故障。消费端group-id按业务分不同业务线用不同消费组避免互相影响offset。auto-offset-reset在核心链路用earliest在通知类场景用latest。消费者方法加try-catch异常消息记录到日志或专用Topic避免因为单条坏消息阻塞整个消费线程。写在最后这套环境从头到尾我实际搭过不下五次踩得最多的坑就是advertised.listeners配置和JsonDeserializer类型绑定。如果你在我给的步骤上跑出了不一样的问题不要慌先把Kafka自带的控制台生产者消费者跑通确认问题出在Kafka还是SpringBoot再逐层排查。能打通控制台消息收发就说明Broker没问题剩下的基本都是配置和序列化的事。按照这个排查思路绝大多数问题都能在十分钟内定位。我个人的习惯是项目里尽量把Kafka的Topic命名、消息格式统一收敛到常量类里发送方和消费方都引用同一套定义避免字符串散落各处导致的大小写不一致问题。刚开始折腾Kafka的时候总是图省事把Topic名称直接写在代码里后面改成常量管理之后维护成本一下子降下来了。希望这篇环境搭建加演示的教程能帮你在SpringBoot集成Kafka的路上少走弯路。
返回列表