本节目标:回答消息可靠性的三个问题——生产端怎么保证不丢、消费端怎么在重复投递下不出错、失败的消息去哪;并给出可落地的 acks、幂等、去重表、重试与死信的写法。
适用版本:Spring Boot 4.1.x(Java 21)
10.3 可靠投递与幂等消费
10.2 把消息发出去了,但「发出去」不等于「不丢、不重」。本节把可靠性的三件事收口:生产端不丢、消费端幂等、失败消息的去处。
10.3.1 三种投递语义在借阅场景下的实际含义
教科书上的三种语义,落到「借阅成功 → 发通知」这条链路上是这样:
| 语义 | 定义 | 借阅场景的后果 | 怎么做到 |
|---|---|---|---|
| 至多一次 | 不重,可能丢 | 通知可能漏发,读者收不到 | 先提交位点再处理;acks=0 |
| 至少一次 | 不丢,可能重 | 通知可能重复,读者收到两封 | 先处理再提交位点;acks=all |
| 恰好一次 | 不丢不重 | 理想状态 | 至少一次 + 幂等消费 |
工程上要记住的一句话:跨系统的「恰好一次」是营销词。 Kafka 的 EOS(exactly-once semantics)只在「Kafka 进、Kafka 出」且全部组件参与同一事务时才成立;一旦链路上有数据库、有邮件网关,就不可能真的端到端恰好一次。真正能落地的组合是「至少一次投递 + 消费端幂等」,把重复挡在业务之外。
所以本节的重点不是追求「恰好一次」,而是:生产端尽量不丢,消费端接受重复并把重复处理干净。
10.3.2 生产者不丢:acks 与幂等
spring.kafka.producer.acks 决定生产者在什么条件下认为发送成功:
| 值 | 含义 | 丢消息风险 |
|---|---|---|
0 | 不等待任何确认 | 高,网络抖动就丢 |
1 | leader 写入即确认 | 中,leader 故障且未同步副本时丢 |
all(等同 -1) | 所有 ISR 副本确认 | 低,生产推荐 |
acks=all 要配合 broker 侧的 min.insync.replicas(最小同步副本数)才有意义:如果 ISR 只剩 1 个副本,acks=all 退化成 acks=1。这两项一个在应用、一个在 broker,需要一起设置。
acks=all 之外还要开生产者幂等,避免「发送超时重试」导致的重复:
spring.kafka.producer.acks=all
spring.kafka.producer.properties.enable.idempotence=true
spring.kafka.producer.properties.max.in.flight.requests.per.connection=5
enable.idempotence=true 让 broker 按「生产者 ID + 序列号」去重,同一生产者重试不会产生重复消息。开启后客户端会强制 acks=all 并保证单分区内有序。
这里要划清边界:生产者幂等只解决「同一个生产者的重试重复」,解决不了「应用层重复调用 send」。业务代码里因为重试逻辑把同一个事件发了两次,生产者幂等拦不住——那要靠消费端幂等。
10.3.3 事务性发送:边界在哪
设置 spring.kafka.producer.transaction-id-prefix 后,自动配置会把 KafkaTemplate 变成事务性的,并注册一个 KafkaTransactionManager bean:
spring.kafka.producer.transaction-id-prefix=loan-tx-
@Bean
KafkaTransactionManager<String, Object> kafkaTransactionManager(
ProducerFactory<String, Object> producerFactory) {
return new KafkaTransactionManager<>(producerFactory);
}
一次事务里发多条消息,要么全成功要么全失败:
kafkaTemplate.executeInTransaction(ops -> {
ops.send("loan-events", loanId.toString(), event);
ops.send("audit-events", loanId.toString(), auditEvent);
return null;
});
真正有价值的是「消费-处理-生产」(read-process-write):把消费者位点提交和生产写入放进同一个 Kafka 事务,实现 Kafka 到 Kafka 的 EOS。监听器里用 sendOffsetsToTransaction 把位点也纳入事务,容器需要配成事务性的。
边界必须说清楚:这套事务只覆盖 Kafka 内部的读写。一旦「处理」这一步是写数据库,数据库不参与 Kafka 事务,就回到了「本地事务与消息投递如何一致」的老问题。此时正确的做法是 10.2 讲的「事务提交后再发消息」加上本节的消费幂等,而不是指望 Kafka 事务包住数据库。
10.3.4 消费者不丢:先落库,再提交位点
消费者侧的可靠性几乎全在「位点提交时机」上。两种顺序,对应两种语义:
| 顺序 | 语义 | 崩溃时的后果 |
|---|---|---|
| 提交位点 → 处理业务 | 至多一次 | 消息丢失 |
| 处理业务 → 提交位点 | 至少一次 | 消息重复 |
生产上要的是后者。手动确认模式(MANUAL)把提交时机交给业务代码:
@KafkaListener(topics = "loan-events", ackMode = "MANUAL")
public void onLoanEvent(LoanEvent event, Acknowledgment ack) {
notificationService.sendBorrowNotice(event); // 1. 先做业务
ack.acknowledge(); // 2. 再提交位点
}
顺序反了就是「至多一次」:ack.acknowledge() 在前,业务还没做,进程崩了,位点已提交,这条消息再也不会被消费。
但「至少一次」意味着重复不可避免:业务处理成功、位点提交前进程崩溃,重启后这条消息会被重放。所以至少一次必然要求消费端幂等。
10.3.5 幂等消费:业务唯一键与去重表
幂等的核心是「同一个事件处理多次,效果等同于一次」。两种落地方式:
- 业务唯一键:让业务本身天然幂等。比如「把借阅状态置为已通知」,重复执行结果一样。
- 去重表:记录已处理过的事件 id,处理前先查。
发通知这类「执行一次就产生一次副作用」的动作,必须用去重表。给通知记录加唯一约束:
CREATE TABLE notification_log (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
event_id VARCHAR(64) NOT NULL,
member_id BIGINT NOT NULL,
channel VARCHAR(16) NOT NULL,
created_at DATETIME NOT NULL,
UNIQUE KEY uk_event_channel (event_id, channel)
);
消费时先写去重记录再发通知,靠数据库唯一约束挡住并发重复:
@KafkaListener(topics = "loan-events", ackMode = "MANUAL")
public void onLoanEvent(LoanEvent event, Acknowledgment ack) {
try {
// 去重记录与业务写入放在同一个本地事务里
notificationLogRepository.insert(event.eventId(), event.memberId(), "EMAIL");
notificationService.sendBorrowNotice(event);
} catch (DataIntegrityViolationException duplicate) {
log.debug("重复事件,跳过 eventId={}", event.eventId());
}
ack.acknowledge();
}
insert 撞上唯一约束会抛 DataIntegrityViolationException,捕获后直接确认——说明另一个实例或上一次重放已经处理过。
关键约束:eventId 必须在生产端生成并全局唯一(用 UUID),不能用消息的 (topic, partition, offset) 做去重键。因为分区重平衡、位点重置、topic 重建都会让 offset 变化,用它去重会误判。
第二个关键约束:去重记录与业务写入要在同一个本地事务里。 如果先写去重记录、事务外再发通知,发通知失败时去重记录已存在,重试会被判为「已处理」而跳过,通知永久丢失。要么两者同事务,要么接受「通知可能重复、绝不漏发」的取舍。
10.3.6 重试与死信:失败消息去哪
业务处理失败(数据库连不上、下游超时)不能无限重试,也不能直接丢。Spring Kafka 提供两条路径。
路径 A:DefaultErrorHandler + 死信发布器(阻塞式重试)
@Bean
DefaultErrorHandler errorHandler(KafkaTemplate<String, Object> template) {
DeadLetterPublishingRecoverer recoverer =
new DeadLetterPublishingRecoverer(template);
ExponentialBackOffWithMaxRetries backOff =
new ExponentialBackOffWithMaxRetries(3); // 最多重试 3 次
backOff.setInitialInterval(1_000L);
backOff.setMultiplier(2.0);
backOff.setMaxInterval(10_000L);
DefaultErrorHandler handler = new DefaultErrorHandler(recoverer, backOff);
handler.addNotRetryableExceptions(IllegalArgumentException.class); // 校验错不重试
handler.setCommitRecovered(true);
return handler;
}
ExponentialBackOffWithMaxRetries 来自 org.springframework.kafka.support(注意不是 org.springframework.util.backoff)。重试在消费线程内阻塞进行,次数用尽后由 DeadLetterPublishingRecoverer 把记录发到死信主题,默认主题名是原 topic 加 .DLT 后缀。
addNotRetryableExceptions(...) 很重要:对「重试也不会成功」的异常(参数校验失败、反序列化错误)直接进死信,否则会白白重试,还拖慢整个分区。
路径 B:@RetryableTopic(非阻塞重试)
把重试拆成一组带延迟的重试 topic,主消费线程不被阻塞:
@RetryableTopic(
attempts = "4",
backOff = @BackOff(delay = 1000, multiplier = 2.0, maxDelay = 10000),
autoCreateTopics = "true",
exclude = {IllegalArgumentException.class})
@KafkaListener(topics = "loan-events")
public void onLoanEvent(LoanEvent event) {
notificationService.sendBorrowNotice(event);
}
@DltHandler
public void handleDlt(LoanEvent event) {
log.error("借阅事件进入死信 loanId={}", event.loanId());
deadLetterService.record(event);
}
启用它需要打开开关(默认关闭):
| 属性 | 默认值 | 说明 |
|---|---|---|
spring.kafka.retry.topic.enabled | false | 必须显式开启 |
spring.kafka.retry.topic.attempts | 3 | 总尝试次数 |
spring.kafka.retry.topic.backoff.delay | 1s | 首次退避 |
spring.kafka.retry.topic.backoff.multiplier | 1 | 退避倍数 |
spring.kafka.retry.topic.backoff.max-delay | 30s | 退避上限 |
spring.kafka.retry.topic.backoff.jitter | 0 | 抖动,避免同时重试 |
注意最后一行的来历:4.0 起 Spring Kafka 的重试能力从 Spring Retry 迁到了 Spring Framework 的 org.springframework.core.retry,原来的 spring.kafka.retry.topic.backoff.random 被更灵活的 jitter 取代。老项目迁移时看到 random 属性要改。
两条路径的取舍:
| 维度 | DefaultErrorHandler | @RetryableTopic |
|---|---|---|
| 是否阻塞 | 阻塞消费线程 | 非阻塞,走独立重试 topic |
| 顺序 | 重试期间阻塞该分区,保序 | 后续消息继续消费,顺序被打乱 |
| 额外资源 | 只需一个死信 topic | 自动创建多个重试 topic |
| 适用 | 重试次数少、要求保序 | 重试间隔长、吞吐优先 |
「借阅通知」对顺序不敏感、又希望失败别拖住整个分区,适合 @RetryableTopic;如果业务强要求分区内保序(比如账户流水),就得用阻塞式重试。
无论走哪条路径,死信主题必须有人消费。死信无人处理等于「消息进了黑洞,日志里还显示成功」——这是最隐蔽的丢消息方式。至少要有一个消费者把死信落库并告警。
10.3.7 常见坑
坑一:先 acknowledge() 再处理业务。 一行顺序之差,语义从「至少一次」滑到「至多一次」,且不报错。
坑二:用 offset 做幂等键。 分区重平衡或位点重置后 offset 会变,同一事件被当成新事件,幂等失效。幂等键必须是生产端生成的业务 id。
坑三:重试无上限、无分类。 毒消息(永远解析不了的报文)会占着分区反复重试,把正常消息堵死。用 addNotRetryableExceptions / exclude 把不可恢复的异常直接送死信。
坑四:去重记录与业务写入不在同一事务。 中间失败会导致「记了去重、没发通知」或反过来,前者丢消息,后者重复。
坑五:把生产者幂等当成端到端幂等。 它只防「同一生产者重试重复」,防不住应用层重复发送,也防不住消费者重放。
坑六:死信主题没人消费。 消息静默进入死信,监控上只看到「消费成功」,故障被掩盖到用户投诉才暴露。
坑七:enable-auto-commit=true 配手动去重。 位点在去重逻辑执行前就提交了,去重表形同虚设。手动 ack 模式必须配 enable-auto-commit=false。
小结
- 跨系统的「恰好一次」不可得,能落地的是「至少一次投递 + 消费端幂等」,把重复挡在业务之外。
- 生产端不丢靠
acks=all+ broker 侧min.insync.replicas+ 生产者幂等;但生产者幂等只防重试重复,防不住应用层重复发送。 - Kafka 事务只覆盖 Kafka 内部的读写;处理步骤一旦写数据库,就得靠「事务提交后发消息 + 消费幂等」而不是指望事务。
- 消费端可靠性的核心是顺序:先处理业务,再提交位点,顺序反了就是至多一次。
- 幂等靠生产端生成的全局唯一
eventId加去重表唯一约束,且去重记录与业务写入必须在同一本地事务。 - 失败消息用
DefaultErrorHandler(阻塞、保序)或@RetryableTopic(非阻塞、走重试 topic)重试,最终进死信主题,且死信必须有消费者与告警。
至此第 10 章把「异步」和「消息」两条线都走完了:进程内用线程池,跨进程用消息队列,两者都要处理上下文与失败。下一章转向数据访问的底层设施——连接池的调优与监控。
阅读导航:上一节:10.2 消息队列集成 · 下一节:11.1 HikariCP 调优 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。