ARTICLE DETAIL

资讯详情

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

Spring Boot集成Kafka实战:从版本选型到消息积压排查

Spring Boot集成Kafka实战:从版本选型到消息积压排查 先说结论Spring Boot 集成 Kafka 本身不算难依赖一加、配置一写、注解一挂就能跑起来。但真正让很多人卡住的是后面的事情——版本选型、监听器机制、消息排错、性能调优、甚至一条 1MB 的消息发不出去的尴尬。这篇文章我不打算写成一个“照着抄就完事”的 demo 合集而是想把我自己从零把 Kafka 整合进 Spring Boot 项目的完整过程、踩过的坑、回头梳理出来的设计思路全部摊开来讲。文章适合这几类人看刚把 Spring Boot 玩熟、想给系统加消息队列的已经在项目里用了 Kafka 但经常遇到消息发不出去、消费延迟、重启后丢数据这类问题的以及准备面试、想把 Kafka 集成这块讲出点深度的。我会从选型思路讲到环境安装从核心配置讲到代码实现最后再把高频问题集中排查一遍全程只讲实际操作里验证过的东西。1. 整体设计思路为什么选 Kafka以及集成前必须想清楚的事1.1 消息队列选型Kafka 到底赢在哪里每次做技术选型我都会先问自己一个问题这个组件解决的是“通信问题”还是“吞吐问题”。如果是简单的服务间同步调用HTTP 和 RPC 足够如果只是要一个“削峰填谷”的中间层RabbitMQ 也很成熟但当业务场景是海量日志、埋点数据、实时链路追踪、或者上游流量波动极大、又要求下游能按自己的节奏消费时Kafka 的高吞吐和日志型存储架构几乎是绕不开的选择。我之前在项目里对比过一组数据同样的三节点集群RabbitMQ 在开启持久化和确认机制后稳定吞吐大概在几千到一万条每秒的量级而 Kafka 只要分区数给够、客户端参数稍微调一调轻松能到几十万条每秒。这不是说 RabbitMQ 不好——它的路由灵活、延迟低、消息确认语义丰富在很多业务场景下更合适。但如果你需要的是“先把大量数据快速收下来再让多个消费者各取所需”Kafka 的发布订阅模型和分区并行机制就是天然的答案。如果非要打一个比方RabbitMQ 像一个服务周到的快递驿站每一件包裹都登记兑换码你凭码领取一码一货Kafka 更像一个巨型仓库里的传送带货品按类别分到不同轨道分区多个工人可以同时在不同的轨道上干活而且传送带本身会把货品存一段时间哪怕工人没空货也不会丢回头再来取就行。这个“传送带自带存储”的特性正是 Kafka 能扛高吞吐、能回溯消费的核心原因。1.2 集成前的核心概念Topic、分区、消费者组在写第一行业务代码之前我建议先把这四个概念在脑海里立起来否则后面配置 KafkaTemplate 和 KafkaListener 时很容易张冠李戴。Topic 是消息的逻辑分类比如“订单事件”“用户登录日志”生产者往 Topic 里写消费者从 Topic 里读。一个 Topic 会被拆成若干个 Partition分区分区是物理上的并行单元也是 Kafka 保证消息有序的最小粒度——同一个分区内的消息是有序的跨分区则不做全局有序承诺。Consumer Group消费者组是 Kafka 实现“一条消息只被组内一个消费者处理”的机制同一个组里的消费者可以均摊分区做到水平扩展。这里有一个新手最容易绕晕的点Kafka 的 Fifo先进先出语义只在分区级别成立。如果你业务上要求某个用户的所有订单事件严格按时间顺序处理那就必须保证这些消息进了同一个分区通常的做法是用“用户ID”作为消息键Kafka 会对同一 Key 哈希到同一分区。这个细节我在项目里吃过亏当时没设 Key消息被均匀打散到多个分区结果同一个订单的创建、支付、完成事件在不同消费者实例上被并行处理出现了状态倒挂。后来统一用 orderId 做 Key问题才消失。还有一个小知识点消费者组里的消费者数量如果大于分区数多出来的消费者会一直空闲因为它没有任何分区可以接管。所以你的消费端扩容并不是无脑加实例而是要看着分区数来定分区不够时优先加分区。1.3 基于幂等设计决定集成方式消息重复是常态在做任何消息系统集成之前团队成员必须先达成一个共识消息重复是必然的不是偶然的。Kafka 的 at-least-once至少一次语义决定了在极端情况下消费者处理完但位移提交失败、网络分区、Rebalance 等可能出现重复消费。所以在设计生产者时把幂等开启在设计消费者时把“按业务唯一键去重”当成默认动作。我在这次集成里铺了这么几条底生产者端开启幂等enable.idempotencetrue消息体里必须携带业务唯一编号比如 eventId / orderId消费者端无论什么场景都先查一遍去重表或者 Redis 判重处理逻辑本身也尽量做成天然幂等比如更新操作直接 set 成最终值而不是“读出来1再写回去”。这套组合拳下来重复消息就算来了也不会污染业务数据。集成方案的技术选型上我坚持在 Spring Boot 项目里通过 spring-kafka 这个官方 Starter 来做而不是自己用原生 Kafka Client 封装。原因有三第一spring-kafka 提供了 KafkaTemplate 和 KafkaListener声明式地搞定生产和消费代码量少一大截第二Spring Boot 的自动配置能帮我们把连接工厂、事务、序列化等基础设施管起来出问题的概率更低第三它在消费端支持批量监听、手动 ACK、异常处理器这些高级能力正好能覆盖真实项目里的绝大多数场景。2. 环境准备Kafka 安装、集群模式与可视化工具2.1 Kafka 本地环境搭建含 Windows 实操Kafka 的安装和历史版本有点纠缠我先说结论新项目优先用 Kafka 3.5 以上版本尽量走 KRaft 模式就是 Kafka 自带的“去 ZooKeeper”的元数据管理方式。因为从 Kafka 3.5 开始ZooKeeper 模式就被标记为弃用到 Kafka 4.0 已经彻底不再支持再抱着 ZK 模式学新版本完全没有必要。Windows 上装 Kafka 其实不复杂我第一次装的时候被网上的老教程绕晕了后来理顺了流程一共就五步装 JDKKafka 3.x 需要 JDK 8 起步建议 JDK 11/17去 Apache 官网下载 Kafka 二进制包比如 kafka_2.13-3.6.2.tgz解压即可打开终端进入 Kafka 目录生成集群唯一 ID.\bin\windows\kafka-storage.bat random-uuid用上一步生成的 ID 格式化存储目录.\bin\windows\kafka-storage.bat format -t uuid -c .\config\server.properties启动 Kafka.\bin\windows\kafka-server-start.bat .\config\server.properties。说到这里要补充一个特别容易踩的坑第 4 步的格式化命令在很多教程里写成kafka-storage.shWindows 上要用kafka-storage.bat另外格式化完成后config 目录下的 log.dirs 指向的目录会生成一堆数据文件这些文件不能手动删除否则启动直接报“存储目录不存在或已损坏”之类的错误。我见过有同事把整个 Kafka 目录拷到别的电脑上忘记重新格式化结果节点起不来——因为 KRaft 元数据里记录了原本的日志目录信息换机器的路径变了自然就崩了。Linux / Mac 上的操作更简单把 .bat 换成 .sh 即可。如果要搭 Kafka 集群只需要准备三台机器或者在一台机器上改三个不同的 server.properties分别设置不同的 listener 端口、日志目录和节点 ID然后用 controller.quorum.voters 把三个节点串起来启动顺序是先逐个格式化、再逐个启动。集群的价值不只是吞吐翻倍更关键的是 Kafka 的副本机制replication.factor能保证单个 Broker 挂掉时数据不丢。2.2 Kafka 可视化工具不要只盯着命令行命令行能完成 Kafka 的所有操作但日常排查问题效率太低。我平时用的可视化工具主要看三类公司运维平台如果已经接入了 Kafka 监控直接在网页上看 Consumer Lag 和 Topic 流量本地开发或者测试环境我用过 Docker 版的 Kafka UI开源项目界面清爽能看 Topic、分区、消费组、消息内容Windows 桌面端还有一个老牌的 Offset Explorer以前叫 Kafka Tool适合不想搞 Docker 的人。说到“kafka有没有ui界面”这个搜索热词我想多说一句Kafka 官方本来就不提供 Web UI生态里的可视化工具不管叫 UI 还是 Manager都是第三方做的。所以你搜到一个界面很花哨的工具不要觉得“这不是官方出品不可靠”——反而应该优先看社区活跃度高的开源项目。实际使用中我推荐在本地开发环境直接跑一个 Kafka UI 的 Docker 容器一条命令就能起来docker run -d --name kafka-ui -p 8080:8080 -e KAFKA_CLUSTERS_0_NAMElocal -e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERSlocalhost:9092 provectuslabs/kafka-ui:latest启动后打开 localhost:8080你就能看到所有 Topic 的分区分布、每个消费组的 Lag堆积量、甚至直接在界面上发送和查看消息。排查问题时这个界面比一条一条敲命令高效太多了。2.3 Kafka Broker 端关键参数集成前先调好底座很多集成的坑根子不在 Spring Boot 代码里而在 Kafka Broker 的默认配置上。有三个参数我必须建议你提前看一眼第一个是 log.retention.hours / log.retention.bytes决定消息在 Broker 上保留多久、多大容量被清理。如果消费者长时间宕机或者消费 Lag 持续堆积消息没有足够保留时间的话回来继续消费时从头开始读可能会发现最早的位移已经被删除了造成“读不到历史消息”的假象。第二个是 auto.create.topics.enable默认是 true生产者和消费者在请求不存在的 Topic 时Broker 会自动创建它。这个特性开发期很友好但生产环境我强烈建议关掉——不然一个拼写错误的 Topic 名Broker 会默默给你建一个新 Topic数据发进去了找不到排查起来极其痛苦。第三个是 message.max.bytes默认是 1048576也就是 1MB。这对应了后台参数里“kafka接收1m”这个高频问题的根源——如果一个 Topic 需要承载超过 1MB 的单条消息光改客户端是没用的Broker 端必须同步调大否则生产端会报“Message size too large”或者超限异常。这个我在第 4 章会展开讲这里先留个锚点。3. Spring Boot 集成实操从依赖到可跑通的生产消费3.1 版本选型Spring Boot、spring-kafka、Kafka Server 三者怎么匹配版本问题是“springboot版本太高”这个热搜词背后最常见的痛点。Spring Boot、spring-kafka、Kafka Server 三者的版本不是随意组合的我把经验整理成一句话看 spring-kafka 的版本它决定了你代码里能用到哪些新特性而 spring-kafka 又和 Spring Boot 的版本强绑定Kafka Server 与客户端之间则是向后兼容的——高版本客户端可以连低版本 Server低版本客户端连高版本 Server 则可能有问题。如果你用的是 Spring Boot 2.7.x对应的 spring-kafka 是 2.8.x此时 Kafka Server 建议至少 2.x 以上如果用的是 Spring Boot 3.xspring-kafka 就是 3.x对应的 Kafka Server 建议 3.4 以上。我当前项目用的是 Spring Boot 3.2.5 spring-kafka 3.2.x Kafka Server 3.6.2这套组合实测稳定批量消费、事务消息、重试机制都能正常工作。这里补一个细节spring-boot-starter 里并不包含 spring-kafka你需要手动引入依赖。Maven 的写法如下dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId /dependencySpring Boot 的依赖版本管理会帮你匹配好 spring-kafka 的版本号所以不需要手动指定 version。需要注意的是如果你引用了 spring-kafka-test记得加scopetest/scope不然打生产包的时候会把测试库带进去体积大且没必要。3.2 配置文件拆解每个参数到底是什么含义很多新手拿到一份 Kafka 配置就照抄抄完出了问题完全不知道从哪个参数下手。所以我把这次项目里用的 application.yml 核心配置贴出来并逐个参数讲清楚为什么这么写spring: kafka: bootstrap-servers: 192.168.1.10:9092,192.168.1.11:9092,192.168.1.12:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer acks: all retries: 3 properties: enable.idempotence: true max.request.size: 5242880 linger.ms: 5 consumer: group-id: order-service-group enable-auto-commit: false key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: earliest properties: max.poll.records: 100 max.partition.fetch.bytes: 1048576 listener: concurrency: 3 missing-topics-fatal: false挑几个关键点解释bootstrap-servers 是集群入口地址多个节点用逗号分隔。消费者和生产者的连接都会从这个列表里的任意一个 Broker 拿到整个集群的元数据所以不需要把每个 Broker 都写全但生产环境建议至少写两个防止单点。acksall 意味着生产者要等所有 ISR同步副本都收到消息后才算发送成功这是“不丢消息”最重要的一道防线。它带来的代价是延迟略高但业务消息值得。enable-auto-commitfalse 这个很多人不理解我多花点笔墨。默认情况下 Spring Kafka 消费端在拉取一批消息后自动提交位移如果此时你的业务逻辑还没处理完进程就挂了位移已经前移重启后会从已提交的位置继续消费那一批消息就丢了。把它关掉改成手动提交才能做到“处理完一条/一批消息后再告诉 Kafka 我读到这了”。手动提交的代码在第 3.3 节里演示。auto-offset-resetearliest 的意思是如果消费者组是新建的、此前没有提交过位移那么从 Topic 最早的消息开始消费。如果设成 latest新组启动后就只消费启动之后产生的新消息之前堆积的数据全都不读。这个参数在你想要“从头补数据”时特别有用但生产环境设成 earliest 一定要小心因为如果消费者组 ID 换了它会把存量数据全部扫一遍可能给下游造成巨大压力。listener.concurrency3 是 KafkaListener 的并发消费线程数。注意它配合的是“每实例并发”并最终受分区数约束——如果你只有一个分区concurrency 设成 10 也没用只会白开 9 个闲置线程。3.3 生产者与消费者代码一个能直接搬到项目里的范式先看生产者。Spring Boot 里注入 KafkaTemplate 就用这段代码我每次写项目都一个套路封装一个只发送业务事件的方法消息体统一转成 JSON 字符串Key 用业务唯一 ID这样能保证同一个业务对象的多次操作落进同一个分区。Service public class OrderEventPublisher { Resource private KafkaTemplateString, String kafkaTemplate; public void publishOrderCreated(OrderCreatedEvent event) { String key String.valueOf(event.getOrderId()); String value JSON.toJSONString(event); // 第二个参数是 Topic 名生产环境建议用常量集中管理 kafkaTemplate.send(order-events, key, value) .whenComplete((result, ex) - { if (ex null) { log.info(消息发送成功, topic{}, partition{}, offset{}, result.getRecordMetadata().topic(), result.getRecordMetadata().partition(), result.getRecordMetadata().offset()); } else { log.error(消息发送失败, key{}, key, ex); } }); } }这里我用的异步回调来确认发送结果。KafkaTemplate.send() 本身是非阻塞的如果你不关心结果直接调用即可但生产环境我强烈建议用 whenComplete 记录一下发送结果——我之前排查过一个问题生产端报“成功”但下游就是收不到结果发现是消息发送到了错误的分区且没有被任何消费者组订阅。如果不做发送结果监控这类问题根本发现不了。再看消费者。KafkaListener 是我的首选方式它比手动写 ConsumerThread 简单太多而且天然支持并发、异常处理、批量消费。这次项目里用的是批量监听手动确认的模式Component public class OrderEventConsumer { KafkaListener(topics order-events, groupId order-service-group) public void onMessage(ListConsumerRecordString, String records, Acknowledgment ack) { try { for (ConsumerRecordString, String record : records) { OrderCreatedEvent event JSON.parseObject(record.value(), OrderCreatedEvent.class); // 消息去重Redis 里 setnx eventId只有第一次处理才返回成功 boolean firstConsumed stringRedisTemplate.opsForValue() .setIfAbsent(event:dedup: event.getEventId(), 1, Duration.ofHours(24)); if (!firstConsumed) { log.info(重复消息跳过, eventId{}, event.getEventId()); continue; } // 处理业务逻辑 orderService.createOrder(event); } // 这批消息全部处理成功后手动提交位移 ack.acknowledge(); } catch (Exception e) { log.error(批量消费失败, size{}, records.size(), e); // 实际项目里这里会结合死信队列或者重试模板先把异常抛出去 throw new KafkaException(消费失败, e); } } }批量监听要在配置文件里加一个监听容器工厂的设置或者在 KafkaListener 上指定 batch 属性另外容器工厂需要设置为 BatchListenerFactoryBean细节可以参考 spring-kafka 官方文档。我这么做最大的好处是吞吐量比单条消费高不少因为一次网络请求拉取一批数据能显著减少网络往返和线程上下文切换。还有一个细节catch 里我直接抛 KafkaException这种处理方式会触发 spring-kafka 的默认重试机制——它会按 ErrorHandler 的配置重新拉取这些消息直到成功或者重试耗尽。如果业务处理失败后不想阻塞整个分区可以改用 SeekToCurrentErrorHandler 之类的错误处理器做更细粒度控制后面第 4 章会聊到。3.4 事务消息与本地事务表说到还是做到的一致性在回答“要不要给 Kafka 加事务”这个问题前我先说一个使用场景用户下单成功后既有数据库订单状态要更新又要发送一条消息给积分服务。如果先更新数据库再发消息消息可能发失败两边数据不一致如果先发消息再更新数据库又可能数据库更新失败但消息已经发出去了。spring-kafka 是支持事务的核心思路是通过 KafkaTransactionManager 把“消息发送”纳入 Spring 管理的事务里。但说实话跨系统的事务一致性在分布式场景里永远有成本我在实际项目里更常采用“本地消息表”这种更朴素的方案先在一个本地数据库事务里写入业务数据消息记录状态为待发送事务提交后再由定时任务或者即时任务把消息发出去发完更新状态。这套方案没有分布式事务的复杂度消息也保证不丢。如果你的业务确实需要 Kafka 事务那么核心步骤如下配置 ProducerFactory 时设置 transactionIdPrefix在发送代码上用 Transactional 包住 kafkaTemplate.send 操作。但必须先想清楚Kafka 的事务只能保证“消息发送”的原子性它管不了你的数据库事务——除非你用 ChainedKafkaTransactionManager 把数据库事务和 Kafka 事务揉在一起而那个东西复杂度极高实际项目里很少用。所以我的建议是先明确你的真实需求再决定是否上事务。4. 生产环境高频问题排查实录4.1 消息延迟高从消费 Lag 到定位瓶颈“kafka消息延迟高”是我看到的热搜词里最典型的生产问题。它的本质是生产速率和消费速率不匹配导致消费组的 Lag积压量越来越大。排查这个问题的第一件事不是看代码而是去看当前消费组的 Lag 数值。如果机器的总消费速率跟不上生产速率第一顺位怀疑的是“分区数不足”。一个分区同一时刻只能被组里的一个消费者实例顺序消费假设单分区处理速度是 5000 条/秒而生产速率是 2 万条/秒这时候你怎么优化消费者代码都没用——瓶颈在分区并行度上。正确做法是给 Topic 增加分区数注意分区数只能增加不能减少同时对应增加消费者实例数或者 listener.concurrency。另一个常见原因是消费线程被“慢业务”阻塞了。比如你拉取了一批消息对每条消息都调用外部 HTTP 接口这个接口平均耗时 500ms那么一次 poll 循环处理 5 分钟都正常。此时看似是 Kafka 的问题实际上是下游依赖的吞吐拖死了消费端。我曾经排查过类似问题最后的结论是消费端在处理消息时同步调了一个第三方风控接口接口超时设置太长导致整个消费者线程池被占满。调短超时、把部分逻辑异步化之后Lag 立刻降了下来。还有一个容易被忽视的配置max.poll.interval.ms默认 300000 毫秒5 分钟。如果消费者两次 poll 之间的最大间隔超过这个值Kafka 会认为消费者已经挂了会对它做主动离组触发 Rebalance。而 Rebalance 期间整个消费组是停止消费的停止期间消息持续堆积恢复后 Lag 又瞬间拉高——形成恶性循环。所以如果你批量处理的任务执行时间可能超过 5 分钟请果断调大这个参数或者把业务拆小。最后给一个排查路径的清单按优先级来先看 Kafka UI 或命令行查 Lag这是第一步定位“到底积压了多少”再看消费者日志确认有没有频繁 Rebalance、超时、连接重置的记录然后看消费端的线程池和下游依赖延迟是否被外部接口拖住最后才考虑调参分区数、fetch 大小、poll 间隔、并发线程数。4.2 单条消息超过 1MB默认限制在哪里怎么调“kafka 接收1m”这个热搜词我猜是“接收 1M 消息”相关的诉求。Kafka 默认的单条消息大小上限是 1MB这是 Broker 端的 message.max.bytes、生产者端的 max.request.size、消费者端的 max.partition.fetch.bytes 三处共同作用的结果。很多人只在 Spring Boot 的 producer 配置里调大了 max.request.size结果还是报错原因就是另外两个位置没调。我当时遇到的一个业务场景是通过消息系统同步一批图片 OCR 识别结果单条消息的 JSON 串因为包含了 base64 的图片内容超过了 2MB。排查链路是这样的先看生产端日志报错“[Record batch is too big]”或者 “Message size too large”再看 Broker 端日志可能没有异常因为 Broker 只是拒绝接收超大消息最后看消费端发现消费者没有任何反应——因为消息根本没进到 Topic 里。那个项目最终的配置改动如下# broker端 message.max.bytes10485760 replica.fetch.max.bytes10485760 # 生产者端spring boot yml spring.kafka.producer.properties.max.request.size10485760 # 消费者端 spring.kafka.consumer.properties.max.partition.fetch.bytes10485760注意 Broker 端的参数是在 server.properties 里改的改完必须重启 Broker 才生效生产端和消费端的参数在 Spring Boot 的 yml 里配置即可不用重启 Broker。这里还有一个更优的实践方案不要用 Kafka 传大块内容。把大内容放到对象存储或者文件服务Kafka 消息里只放文件 ID 和元数据消费者再按 ID 去拉取。这样做的好处是 Kafka 集群不用为了极少数的“大消息”调整全局参数——毕竟 message.max.bytes 一旦调大会影响集群所有 Topic 的吞吐和内存占用得不偿失。4.3 版本太高导致的各种不可名状问题“springboot版本太高”这个热搜词背后其实是两种常见情况一种是项目引入了最新版 Spring Boot但 spring-kafka 没跟上接口报错另一种是本地开发环境连的 Kafka Server 版本太老新的客户端协议不兼容。我先说第一种。Spring Boot 3.x 和 2.x 在 Kafka 集成上有一些破坏性变化比如 spring.factories 变成了 AutoConfiguration.imports、部分 KafkaProperties 的属性路径有调整。如果你用 Spring Boot 3.2 但引入的是一个很老的 spring-kafka 2.3 版本编译都可能过不去。解决办法很简单不要手动指定 spring-kafka 的 version交给 Spring Boot 的 dependency management 去管。Boot 3.2.5 会对应引入 3.1.x 的 spring-kafka这个版本对 Kotlin 和 Jackson 的处理都比较正常。第二种情况更隐蔽。我踩过的一个坑是本地 Kafka Server 是 2.4 版本但项目里的 spring-kafka 是 3.x启动时报 “Unsupported version” 之类的异常。查了资料才发现Kafka 客户端从 3.0 开始就放弃了对 2.x 老版本的兼容——客户端只能连比自己版本低不超过一定范围的 Server。所以如果你公司核心集群很老业务代码的 Spring Boot 版本不能升太高否则 Kafka 客户端一升级就连不上了。反过来如果集群是新的客户端老一点反而问题不大。排查版本问题的通用方法是看启动日志里的 Kafka 版本协商信息或者直接用官方自带脚本测试连通性kafka-broker-api-versions.sh --bootstrap-server localhost:9092这个命令能直接输出 Broker 支持的 API 版本列表把它和客户端日志报错对比一下问题基本就清楚了。4.4 消费者进程重启后丢消息手动确认与 Rebalance 的博弈我遇到过好多次类似反馈“服务重启后之前明明收到过的消息又消费了一遍或者不消费了”。前者是重复消费后者是丢消息。这两个问题在 Kafka 消费端特别常见根子都在位移提交的时机上。先说丢消息。如果开了自动提交enable-auto-committrue且提交间隔较短可能出现“消息已经被拉取到本地、位移已提交、但业务代码还没处理完”的场景——此时如果进程重启这批消息不会再被拉取相当于丢了。解决办法就是前面说的改成手动提交业务处理成功后再 ack。但手动提交也有两个分支同步提交ack.acknowledge() 阻塞等待结果和异步提交ack.acknowledge() 回调。我习惯用同步提交代码逻辑更直观虽然吞吐略低如果追求极限吞吐可以异步提交但一定要在回调里记录失败日志因为异步提交失败时不会通知到你静默丢位移的风险非常大。再说重复消费。如果你在 ack 之前业务代码抛了异常spring-kafka 默认会让这批消息进入重试重试耗尽后可能会再次投递这就是重复消费的来源。解决重复消费的方法我在第一章就提过业务侧做幂等。我现在所在团队有一条铁律“任何消费端逻辑必须天然幂等或者有去重机制否则不允许接入 Kafka。”这条规矩帮我挡住了很多线上事故。Rebalance 和位移提交还有个交互细节值得一说当一个消费者实例退出或者加入消费组Kafka 会触发 Rebalance此时未提交的位移会回滚到上次 commit 的位置。如果你的消费者正在处理一批消息刚处理完还没来得及 ack正好发生 Rebalance那这批消息在 Rebalance 结束后会被重新拉取一遍业务逻辑会再执行一遍。这个场景是重复消费的经典入口也是我认为“幂等是硬需求”的重要原因——你没有办法保证 100% 避开这种情况只能保证重复执行的结果是一样的。5. 集成后的扩展思路Kafka 集成不是把消息发出去收回来就结束了真正有复杂度的是围绕它搭建的生态能力。这次项目完成后我又在上游加了消息轨迹日志每条消息都有 EventId、Time、Topic、Partition、Offset 的完整链路线索在下游加了消费失败的死信队列处理失败的消息投递到专门的“下划线开头”的 DLT Topic方便人工排查还在中间引入了基于 KafkaListener 的指数退避重试——第一次失败 1 秒后重试第二次 2 秒第三次 4 秒最多 5 次。这一套下来Kafka 从“能用”变成了“好用”日常运维和排障效率提升非常明显。最后分享一个个人习惯新项目里只要涉及 Kafka 配置我会把所有 Topic 名集中到一个常量类里所有消费者 Group ID 用“业务名-环境-用途”的命名规范比如 order-service-dev-consumer所有消息体必须包含 eventId 和 timestamp 字段。这些约定在项目初期看似繁琐等系统上线、每天千万级消息在集群里流转的时候你就会发现这些“小规矩”就是排查一线问题时的指南针。这次集成前后花了我大概三天时间其中两天在和版本、配置缠斗真正写业务代码只有半天——但正是这些藏在配置和机制背后的细节决定了你的消息系统是稳定支撑业务还是变成一个随时可能爆雷的定时炸弹。
返回列表