《Spring Boot 实战》10.3 可靠投递与幂等消费

把消息投递做可靠:讲清生产者 acks 与幂等、事务性发送的真实边界,消费者手动 ack 与「先落库再提交位点」的先后顺序、业务唯一键去重、重试与死信主题的落地写法,并用一张表对比至少一次、至多一次、恰好一次在借阅通知场景下的真实含义与取舍。

本节目标:回答消息可靠性的三个问题——生产端怎么保证不丢、消费端怎么在重复投递下不出错、失败的消息去哪;并给出可落地的 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不等待任何确认高,网络抖动就丢
1leader 写入即确认中,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.enabledfalse必须显式开启
spring.kafka.retry.topic.attempts3总尝试次数
spring.kafka.retry.topic.backoff.delay1s首次退避
spring.kafka.retry.topic.backoff.multiplier1退避倍数
spring.kafka.retry.topic.backoff.max-delay30s退避上限
spring.kafka.retry.topic.backoff.jitter0抖动,避免同时重试

注意最后一行的来历: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 调优 。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「java」更多文章

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