本节目标:把「借阅成功后发事件通知」从进程内异步升级为跨服务消息投递,掌握 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 / JacksonJsonDeserializer | Jackson 3(tools.jackson.databind.json.JsonMapper) | 4.1 首选 |
JsonSerializer / JsonDeserializer | Jackson 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_IMMEDIATE | acknowledge() 立即提交 | 精确控制,有同步开销 |
默认的 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 与发送回调。
选型对比:
| 维度 | Kafka | RabbitMQ |
|---|---|---|
| 模型 | 分布式日志,消费方维护位点 | 队列,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 可靠投递与幂等消费 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。