《Spring Boot 实战》10.2 消息队列集成

以 Spring Kafka 为主、RabbitMQ 为对照,讲清借阅事件通知的生产端与消费端落地:starter 与模块化依赖、序列化选型、@KafkaListener 容器工厂与并发度、消费位点提交模式、批量消费,以及生产端发送回调与分区顺序保证。

本节目标:把「借阅成功后发事件通知」从进程内异步升级为跨服务消息投递,掌握 Spring Kafka 的生产端与消费端配置、位点提交模式与批量消费,并对照 RabbitMQ 说清选型依据。
适用版本:Spring Boot 4.1.x(Java 21)

10.2 消息队列集成

10.1 的 @Async 只解决了「当前进程内不阻塞」。一旦通知要跨服务、要保证不丢、要被多个下游各自消费,就得把动作投递到消息中间件。本节用 Spring Kafka 落地「借阅成功 → 发通知事件」,并以 RabbitMQ 作为对照。

10.2.1 什么时候该引入消息,什么时候不该

消息中间件不是「异步」的同义词。先把三种手段的边界划清:

手段边界典型场景
直接 RPC强一致、要立即拿到结果借书时校验库存
@Async同进程、可丢、无重试诉求记录访问日志
消息队列跨进程、要解耦、要削峰、要可重放借阅成功后通知下游

引入消息换来三样东西:解耦(生产方不需要知道谁消费)、削峰(突发流量进队列慢慢消费)、可重放(Kafka 保留日志,可以重放历史事件)。

代价同样明确:系统从「一次调用」变成了「最终一致」。借阅接口返回成功时,通知可能还没发出去;下游可能重复收到同一条事件;消息可能延迟甚至丢失(取决于配置)。凡是「必须立刻一致」的动作——比如扣减库存——不要放到消息里做,留在事务里。

10.2.2 依赖与 starter

Spring Boot 4.x 做了模块化重构:功能模块叫 spring-boot-<technology>,starter 叫 spring-boot-starter-<technology>。Kafka 的 starter 是:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-kafka</artifactId>
</dependency>

测试支持单独一个 starter:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-kafka-test</artifactId>
    <scope>test</scope>
</dependency>

对应的自动配置模块是 spring-boot-kafka,包名 org.springframework.boot.kafka.autoconfigure。不要再写老式的 org.springframework.kafka:spring-kafka 裸依赖——starter 会把它和 kafka-clients 一起带进来。4.1.1 实测这条版本线是:Spring Kafka 4.1.1、kafka-clients 4.2.1。

RabbitMQ 的对照 starter 是 spring-boot-starter-amqp(模块 spring-boot-amqp,包 org.springframework.boot.amqp.autoconfigure,Spring AMQP 4.1.1)。两者可以同时存在,互不冲突。

10.2.3 生产端:KafkaTemplate 与序列化

Spring Kafka 自动配置会提供一个 KafkaTemplate<K, V>。发消息前先定序列化——这里有个 4.x 特有的坑。

Spring Boot 4.1 默认用 Jackson 3(包名 tools.jackson),而 Spring Kafka 里有两套 JSON 序列化器:

类底层4.x 用法
JacksonJsonSerializer / JacksonJsonDeserializerJackson 3(tools.jackson.databind.json.JsonMapper)4.1 首选
JsonSerializer / JsonDeserializerJackson 2(com.fasterxml.jackson)旧写法,仍可用

沿用旧的 JsonSerializer 会让 Kafka 侧的 JSON 与 Web 层的 JSON 走两套 mapper,时间格式、命名策略可能不一致。4.1 应该用 JacksonJsonSerializer 系列。

配置生产端:

spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JacksonJsonSerializer
spring.kafka.producer.acks=all
spring.kafka.template.default-topic=loan-events

发送时 KafkaTemplate.send(...) 返回 CompletableFuture<SendResult<K, V>>,用它挂回调,不要用 try/catch 包住 send 就以为捕获到了发送失败——send 本身是异步的,网络层的失败在 future 里:

public void publish(LoanEvent event) {
    kafkaTemplate.send("loan-events", event.loanId().toString(), event)
        .whenComplete((result, ex) -> {
            if (ex != null) {
                log.error("借阅事件发送失败 loanId={}", event.loanId(), ex);
            } else {
                log.debug("借阅事件已发送 partition={} offset={}",
                    result.getRecordMetadata().partition(),
                    result.getRecordMetadata().offset());
            }
        });
}

把 key 设成 loanId 是有意的:Kafka 按 key 分区,同一借阅的事件永远落到同一分区,下游按借阅聚合时就天然有序。如果 key 传 null,事件会轮询到各分区,顺序无法保证。

10.2.4 消费端:@KafkaListener 与容器工厂

消费端用 @KafkaListener,容器工厂由自动配置提供(ConcurrentKafkaListenerContainerFactory)。

spring.kafka.consumer.group-id=loan-notifier
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.enable-auto-commit=false
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JacksonJsonDeserializer
spring.kafka.listener.concurrency=3
@Component
public class LoanEventConsumer {

    @KafkaListener(topics = "loan-events", groupId = "loan-notifier")
    public void onLoanEvent(LoanEvent event) {
        notificationService.sendBorrowNotice(event);
    }
}

几个必须理解的点:

  • concurrency 是消费者实例数,不是「线程数」的泛称。设成 3 意味着这个监听器起 3 个消费者线程,分区数要 ≥ 并发数才有意义,否则多出来的消费者会空转。
  • enable-auto-commit=false 是关键。自动提交由消费者定时在后台提交位点,可能在消息真正处理完之前就提交了——一旦处理失败,位点已经前进,这条消息就丢了。生产上关掉它,交给容器按确认模式提交(10.3 展开)。
  • auto-offset-reset 只在「该消费组没有已提交位点」时生效。设 earliest 会从头消费,新接入的组要注意别把历史全量灌进来。

反序列化目标类型怎么定?@KafkaListener 的方法签名里写了 LoanEvent,但 JacksonJsonDeserializer 默认还会看消息头里的类型信息。跨服务场景下类型头不可信,应当显式关闭并锁定目标类型:

@Bean
ConsumerFactory<String, LoanEvent> loanEventConsumerFactory() {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    JacksonJsonDeserializer<LoanEvent> valueDeser =
        new JacksonJsonDeserializer<>(LoanEvent.class);
    valueDeser.setUseTypeHeaders(false);   // 不信任消息头里的类型
    return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), valueDeser);
}

生产端对应地关掉类型信息写入:JacksonJsonSerializer.setAddTypeInfo(false)。两边都不写类型头,靠代码里的 LoanEvent.class 定类型,既安全又省掉一处版本耦合。

10.2.5 消费位点提交模式

位点(offset)提交的时机决定了「至少一次」还是「至多一次」。容器用 ContainerProperties.AckMode 控制,可用 spring.kafka.listener.ack-mode 全局设,也可在 @KafkaListener(ackMode = "...") 上单设:

AckMode提交时机语义
RECORD每条记录处理完就提交接近至少一次,开销大
BATCH(默认)一批 poll 回来的记录全部处理完提交至少一次,吞吐好
TIME按时间间隔提交吞吐高,失败窗口大
COUNT按记录数提交同上
COUNT_TIME时间或数量任一满足同上
MANUAL业务代码调用 acknowledge() 后,在下一次 poll 时提交精确控制
MANUAL_IMMEDIATEacknowledge() 立即提交精确控制,有同步开销

默认的 BATCH 就够大多数场景用。只有当你需要「处理成功后才提交」的精确控制时,才切到 MANUAL/MANUAL_IMMEDIATE,并在方法里注入 Acknowledgment:

@KafkaListener(topics = "loan-events", ackMode = "MANUAL")
public void onLoanEvent(LoanEvent event, Acknowledgment ack) {
    notificationService.sendBorrowNotice(event);  // 先做业务
    ack.acknowledge();                            // 再提交位点
}

顺序不能反。先 acknowledge() 再处理业务,等于把 10.3 要讲的「至多一次」问题手动造了出来。

10.2.6 批量消费

单条消费在每条记录上都要走一次反序列化与提交,吞吐受限。需要时切成批量模式:

spring.kafka.listener.type=batch
spring.kafka.listener.ack-mode=MANUAL
@KafkaListener(topics = "loan-events", ackMode = "MANUAL")
public void onBatch(List<LoanEvent> events, Acknowledgment ack) {
    notificationService.sendBatch(events);
    ack.acknowledge();
}

spring.kafka.listener.type 的默认值是 single,改成 batch 后方法签名要收 List<...>。批量模式与手动 ack 搭配时要注意:整批一起提交位点,意味着批内有一条失败,其余成功的也会被重放——所以批量消费必须配幂等(10.3),否则重复通知会成倍放大。

10.2.7 完整切片:借阅成功发事件

把上面的碎片拼成一条能跑的链路。先定义事件(用 record,天然不可变、天然可序列化):

public record LoanEvent(
        String eventId,
        Long loanId,
        Long bookId,
        Long memberId,
        String type,
        Instant occurredAt) {
}

借阅服务在事务提交后发事件。关键顺序:先提交数据库事务,再发消息,否则可能出现「消息发出去了、事务回滚了」的幽灵事件:

@Service
public class LoanService {

    private final LoanRepository repository;
    private final LoanEventPublisher publisher;

    @Transactional
    public Loan borrow(Long bookId, Long memberId) {
        Loan loan = repository.save(Loan.create(bookId, memberId));
        // 事务内只登记事件,不直接发送
        return loan;
    }
}

更稳妥的做法是用 Spring 的 @TransactionalEventListener,让发送动作绑定在事务提交之后:

@Component
public class LoanEventListener {

    private final LoanEventPublisher publisher;

    @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
    public void onCommitted(LoanBorrowedDomainEvent event) {
        publisher.publish(LoanEventMapper.toMessage(event));
    }
}

消费端把事件转成通知,并保留重放安全(10.3 会补上幂等):

@Component
public class LoanEventConsumer {

    private final NotificationService notificationService;

    @KafkaListener(topics = "loan-events", groupId = "loan-notifier")
    public void onLoanEvent(LoanEvent event) {
        log.info("收到借阅事件 eventId={} loanId={}", event.eventId(), event.loanId());
        notificationService.sendBorrowNotice(event);
    }
}

启动后日志里能看到容器起来的证据,下面是示例输出(本机未接入真实 Kafka,仅示意格式):

INFO --- [ntainer#0-0-C-1] o.s.kafka.listener.KafkaMessageListenerContainer : loan-notifier: partitions assigned: [loan-events-0, loan-events-1, loan-events-2]
INFO --- [ntainer#0-0-C-1] c.e.loan.LoanEventConsumer : 收到借阅事件 eventId=8f2c loanId=42

partitions assigned 这行是排查「消费者为什么没收到消息」的第一现场:如果分区没分配,多半是 group-id 写错或分区被别的实例抢走。

10.2.8 RabbitMQ 对照

RabbitMQ 与 Kafka 的模型不同:Kafka 是「日志 + 位点」,RabbitMQ 是「队列 + 投递确认」。同样的通知链路,RabbitMQ 版本是:

@Component
public class LoanNoticeListener {

    @RabbitListener(queues = "loan.notice.queue")
    public void onNotice(LoanEvent event) {
        notificationService.sendBorrowNotice(event);
    }
}

对应配置:

spring.rabbitmq.host=localhost
spring.rabbitmq.listener.simple.acknowledge-mode=manual
spring.rabbitmq.publisher-confirm-type=correlated
spring.rabbitmq.publisher-returns=true

publisher-confirm-type=correlated 让生产端能拿到每条消息的确认,publisher-returns=true 处理「路由不到队列」的返回。这两项对应 Kafka 侧的 acks 与发送回调。

选型对比:

维度KafkaRabbitMQ
模型分布式日志,消费方维护位点队列,broker 维护投递状态
顺序分区内有序队列内有序
重放原生支持,改位点即可重放不支持,消费即出队
吞吐量级高(顺序写日志)中(按队列/连接数)
路由弱,靠 topic 命名约定强,exchange 支持多种路由
典型用途事件流、日志、CDC任务队列、复杂路由

一句话:要「事件可重放、高吞吐」选 Kafka;要「复杂路由、每消息精确投递」选 RabbitMQ。 「借阅事件通知」偏事件流,用 Kafka 更合适;如果是「生成一份 PDF 报表」这类一次性任务队列,RabbitMQ 更省心。

10.2.9 常见坑

坑一:在事务里直接 kafkaTemplate.send。 事务回滚了消息却已发出,产生幽灵事件。用 @TransactionalEventListener(AFTER_COMMIT) 把发送挪到提交之后。

坑二:生产端以为 send 同步抛异常。 send 返回 future,失败在回调里。不挂 whenComplete 就等于静默丢弃发送失败。

坑三:enable-auto-commit=true。 位点在消息处理完之前就前进,处理失败即丢消息。生产一律关掉。

坑四:concurrency 大于分区数。 多出来的消费者分不到分区,空占资源。并发数应 ≤ 分区数。

坑五:4.1 里继续用 Jackson 2 的 JsonSerializer。 与 Web 层 Jackson 3 的 mapper 不一致,出现时间格式、字段命名的偏差。改用 JacksonJsonSerializer/JacksonJsonDeserializer。

坑六:把「消息发出」等同于「下游已处理」。 两者之间隔着网络、队列和消费位点,本节的链路只保证「投递」,真正的可靠性与幂等是下一节的内容。

小结

  • 引入消息中间件的收益是解耦、削峰、可重放,代价是从「一次调用」变成「最终一致」;必须立即一致的动作不要放进消息。
  • 4.x 模块化后 Kafka 用 spring-boot-starter-kafka(模块 spring-boot-kafka),RabbitMQ 用 spring-boot-starter-amqp。
  • 4.1 默认 Jackson 3,Kafka JSON 序列化应用 JacksonJsonSerializer/JacksonJsonDeserializer,并关闭类型头、锁定目标类型。
  • KafkaTemplate.send 异步返回 CompletableFuture,失败只在回调里;用 key 保证同一借阅的事件分区内有序。
  • 消费端 concurrency 是消费者实例数且不应超过分区数;enable-auto-commit=false,位点提交交给 AckMode。
  • 事务与消息的顺序是「先提交事务,再发消息」,用 @TransactionalEventListener(AFTER_COMMIT) 保证。

消息发出去了,但「发出去」不等于「不丢、不重」。生产者侧怎么保证不丢、消费者侧怎么保证重复投递下不出错,是 10.3 要收的口。

阅读导航:上一节:10.1 @Async 与线程池 · 下一节:10.3 可靠投递与幂等消费 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java」更多文章

  1. 《Spring Boot 入门》18.3 打包与运行
  2. 《Spring Boot 入门》18.2 实现
  3. 《Spring Boot 入门》18.1 需求与设计