Modern Architecture
& Coding Solutions

Spring Boot + Kafka 深度实战:可靠消息传递与 Exactly-Once 语义

从消息丢失到精确一次处理,一篇讲透生产级 Kafka 的消息可靠性工程

先讲一件真事。

去年某个电商大促的凌晨,我收到一条报警——支付成功后的订单确认消息没有推送给用户。排查了半天,发现是 Kafka 生产者发送消息时网络超时,重试机制把一条消息写了两遍,下游的积分服务收到两条重复消息,给同一个订单加了两次积分。虽然最后人工修复了数据,但那个凌晨的补救工作让我记忆深刻。

那次之后我开始认真琢磨一个问题:Kafka 号称“可靠消息队列”,但在真实的生产环境里,消息丢失和重复消费几乎每天都在发生。 区别只是你有没有发现而已。

Kafka 给了你三种投递语义的选择——最多一次、至少一次、精确一次。但很多人以为“精确一次”是默认行为,或者觉得配置了 enable.idempotence=true 就万事大吉。实际上,精确一次语义的边界比你想象的要窄得多,而代价也比想象中大得多。

本文会从 Spring Boot 3.5 集成 Kafka 开始,逐层拆解消息可靠性的各个环节——生产端怎么配置才能不丢消息、不重复发送;消费端怎么处理才能保证每条消息都被正确处理且只处理一次;以及那个被问过无数次的问题:Exactly-Once 到底能不能做到?

一、三种投递语义:先搞清楚你真正需要什么

在开始写代码之前,先花几分钟搞清楚这三个概念。很多线上事故,根源就在于选了错误的语义。

1.1 At-most-once(最多一次)

生产者发完消息就不管了,不关心 Broker 有没有收到。如果网络抖动或者 Broker 暂时不可用,消息就丢了。

适用场景:日志采集、监控指标上报这类允许少量丢失的数据。

1.2 At-least-once(至少一次)

生产者会重试,直到 Broker 确认收到为止。但问题在于——如果 Broker 已经写入了消息,只是网络超时导致确认包没回来,生产者重试就会把同一条消息再写一遍,造成重复。

这是 Kafka 的默认行为,也是大多数生产环境实际在用的模式

适用场景:大部分业务场景——订单、支付、库存等。前提是下游必须能处理重复消息(幂等)。

1.3 Exactly-once(精确一次)

消息在 Broker 中只被写入一次,且只被消费一次。

但这里有个很多人忽略的细节:Kafka 的精确一次语义只覆盖“Kafka 到 Kafka”的路径。也就是说,它保证的是生产者写到 Topic 不重复、Kafka Streams 消费再产出到另一个 Topic 不重复。它不保证你的数据库写入、外部 API 调用、发短信这些操作也精确一次。

如果你的业务逻辑里调了第三方支付接口,Kafka 的事务回滚管不了那个接口——消息重跑一遍,支付就可能扣两次。

搞清楚这三个语义的边界,比学会配置参数更重要。

二、生产端可靠性:别让消息“发丢了”

消息从生产端到 Broker,这一路有多个环节可能出问题。我们先从生产端的配置开始。

2.1 三个最关键的参数

Spring Boot 的 Kafka 自动配置已经帮我们做了很多事情,但默认值通常不够可靠。下面这三个参数是生产端可靠性的基石:

acks:控制 Broker 在什么情况下确认消息写入成功。

含义可靠性
0不等待任何确认最高性能,消息可能丢
1Leader 副本写入即确认折中方案,Leader 宕机可能丢
all(或 -1ISR 中所有副本都写入才确认最高可靠性

retries:发送失败后的重试次数。生产环境建议设成一个很大的值(如 Integer.MAX_VALUE),让重试机制配合幂等性来保证最终成功。

enable.idempotence:幂等生产者开关。开启后,Broker 会给每个生产者分配一个 PID(Producer ID),并为每个分区维护一个序列号。如果生产者重试一条已经写成功的消息,Broker 看到相同的序列号就直接丢弃,不会重复写入。

注意:开启幂等性后,acks 会被自动设为 allretries 会被设为 Integer.MAX_VALUEmax.in.flight.requests.per.connection 会被设为 5(或 1,取决于 Kafka 版本)。你不需要再手动配置这几个参数。

2.2 生产端配置示例

spring:
  kafka:
    producer:
      bootstrap-servers: localhost:9092
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      # 可靠性核心配置
      properties:
        enable.idempotence: true          # 开启幂等,防止重复写入
        acks: all                         # 所有副本确认才算成功
        retries: 2147483647               # 无限重试(配合幂等)
        max.in.flight.requests.per.connection: 5
        delivery.timeout.ms: 120000       # 整个发送流程的超时
        request.timeout.ms: 30000         # 单次请求超时

delivery.timeout.ms 是一个容易被忽略的参数。它规定了从消息被 send() 到收到 Broker 确认或抛出异常的总时间上限。如果超过了这个时间,即使重试还没结束,生产者也会放弃并抛出异常。

2.3 发送消息时的异常处理

KafkaTemplatesend() 方法是异步的,返回 ListenableFuture<SendResult>。生产环境里,同步发送 + 捕获异常是最简单的可靠性保障方式:

@Service
@Slf4j
public class OrderEventProducer {

    private final KafkaTemplate<String, Object> kafkaTemplate;

    public OrderEventProducer(KafkaTemplate<String, Object> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    /**
     * 同步发送订单事件
     * 只有 Broker 确认写入成功后才返回,否则抛出异常
     */
    public void sendOrderEvent(OrderEvent event) {
        String topic = "order-events";
        String key = event.getOrderId();  // 用订单ID做Key,保证同一订单的消息有序

        try {
            SendResult<String, Object> result = kafkaTemplate
                .send(topic, key, event)
                .get(10, TimeUnit.SECONDS);  // 同步等待,设置超时

            RecordMetadata metadata = result.getRecordMetadata();
            log.info("消息发送成功: topic={}, partition={}, offset={}, orderId={}",
                metadata.topic(), metadata.partition(), metadata.offset(),
                event.getOrderId());

        } catch (TimeoutException e) {
            // 超时异常——可能已经写入,也可能没写入
            // 此时无法确定消息状态,需要配合幂等性处理
            log.error("消息发送超时: orderId={}", event.getOrderId(), e);
            throw new RuntimeException("消息发送超时", e);

        } catch (Exception e) {
            log.error("消息发送失败: orderId={}", event.getOrderId(), e);
            throw new RuntimeException("消息发送失败", e);
        }
    }
}

同步发送会阻塞线程,对吞吐量有影响。如果追求高吞吐,可以用异步发送 + 回调:

kafkaTemplate.send(topic, key, event)
    .addCallback(
        result -> log.info("发送成功: {}", result.getRecordMetadata().offset()),
        failure -> log.error("发送失败", failure)
    );

异步发送不阻塞主线程,但异常需要在回调里处理。两种方式各有取舍,看业务对可靠性的要求。

三、事务:让多条消息“要么全发,要么全不发”

幂等生产者解决了“单条消息不重复”的问题,但它解决不了“多条消息一起发,其中一条失败怎么办”的问题。

举个例子:订单服务创建一笔订单,需要同时发两条消息——一条给库存服务(扣库存),一条给积分服务(加积分)。如果库存消息发出去了,积分消息发送失败,系统就处于不一致的状态。

Kafka 事务解决了这个问题:一批消息要么全部成功写入,要么全部失败回滚

3.1 开启事务

spring:
  kafka:
    producer:
      transaction-id-prefix: tx-order-   # 事务ID前缀,必填
      # 开启事务会自动启用幂等性,所以不用单独配置 enable.idempotence

配置了 transaction-id-prefix 后,Spring Boot 会自动配置 KafkaTransactionManager

3.2 事务性发送

@Service
@Slf4j
public class OrderEventProducer {

    private final KafkaTemplate<String, Object> kafkaTemplate;

    public OrderEventProducer(KafkaTemplate<String, Object> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    /**
     * 事务性发送:订单事件 + 库存扣减事件 要么都成功,要么都失败
     */
    public void sendOrderEventsWithTransaction(OrderEvent orderEvent, 
                                                StockEvent stockEvent) {
        // 在事务中执行
        kafkaTemplate.executeInTransaction(operations -> {
            // 发送订单事件
            operations.send("order-events", 
                orderEvent.getOrderId(), orderEvent);
            // 发送库存扣减事件
            operations.send("stock-events", 
                stockEvent.getOrderId(), stockEvent);

            log.info("事务性消息发送完成: orderId={}", orderEvent.getOrderId());
            // 返回 true 表示提交事务
            return true;
        });
    }
}

executeInTransaction() 方法保证:如果两个 send() 都成功,事务提交;如果任何一个失败,整个事务回滚。

3.3 事务的代价

事务不是免费的。开启事务后,Broker 需要额外的事务协调器(Transaction Coordinator)来管理事务状态,使用两阶段提交协议来保证原子性。

实测数据显示:开启事务后,端到端平均延迟从 25ms 上升到 45ms,TPS 降低约 15%

所以,只在真正需要原子性的场景使用事务——比如订单创建 + 库存扣减必须同时成功或同时失败。如果只是发一条普通日志消息,完全没必要开事务。

四、消费端可靠性:别让消息“白处理了”

生产端把消息可靠地送到了 Broker,消费端如果处理不好,前面的努力就白费了。

4.1 关闭自动提交,手动控制 Offset

Spring Kafka 默认是自动提交 Offset 的。这意味着消费者拉取到消息后,不管业务逻辑是否处理成功,Offset 都会往前推进。

如果业务处理失败了,Offset 已经提交了,这条消息就永远“消失”了——消息丢失

正确的做法是:关闭自动提交,业务处理成功后再手动提交 Offset

spring:
  kafka:
    consumer:
      enable-auto-commit: false  # 关闭自动提交
      # 手动提交模式

手动提交有两种方式:

方式一:在 @KafkaListener 中注入 Acknowledgment

@Component
@Slf4j
public class OrderEventConsumer {

    private final OrderService orderService;

    public OrderEventConsumer(OrderService orderService) {
        this.orderService = orderService;
    }

    @KafkaListener(topics = "order-events", groupId = "order-consumer-group")
    public void consume(ConsumerRecord<String, OrderEvent> record, 
                        Acknowledgment ack) {
        try {
            OrderEvent event = record.value();
            log.info("收到订单事件: orderId={}, status={}", 
                event.getOrderId(), event.getStatus());

            // 处理业务逻辑
            orderService.processOrderEvent(event);

            // 业务处理成功,手动提交 Offset
            ack.acknowledge();
            log.info("消息处理成功,Offset 已提交: partition={}, offset={}", 
                record.partition(), record.offset());

        } catch (Exception e) {
            log.error("消息处理失败: orderId={}", 
                record.value().getOrderId(), e);
            // 不提交 Offset——消息会被重新消费
            // 注意:这里如果一直失败,会导致无限重试,需要配合重试策略
        }
    }
}

方式二:使用 AckMode 配置

spring:
  kafka:
    listener:
      ack-mode: MANUAL_IMMEDIATE  # 手动立即提交

MANUAL_IMMEDIATE 模式下,调用 ack.acknowledge() 会立即提交 Offset。

4.2 幂等消费:防止重复处理

手动提交 Offset 解决了“消息丢失”的问题,但引入了另一个问题——重复消费

如果业务处理成功了,但在提交 Offset 之前发生了异常(比如网络抖动、应用重启),消费者重启后会重新拉取这条消息,业务逻辑会再执行一遍。

解决方案:让消费逻辑具有幂等性

幂等的意思是:同一条消息处理多次,效果和只处理一次相同

常见实现方式:

方式一:数据库唯一约束

@Service
@Slf4j
public class OrderEventConsumer {

    private final OrderRepository orderRepository;

    @KafkaListener(topics = "order-events", groupId = "order-consumer-group")
    public void consume(ConsumerRecord<String, OrderEvent> record, 
                        Acknowledgment ack) {
        OrderEvent event = record.value();
        try {
            // 使用订单ID作为唯一键,数据库层面保证幂等
            // 如果订单已存在,INSERT 会抛出 DuplicateKeyException
            orderRepository.insertOrder(event);
            ack.acknowledge();
        } catch (DuplicateKeyException e) {
            // 订单已存在,说明是重复消息,直接确认
            log.warn("订单已存在,跳过重复消息: orderId={}", event.getOrderId());
            ack.acknowledge();
        } catch (Exception e) {
            log.error("消息处理失败: orderId={}", event.getOrderId(), e);
            // 不提交 Offset,等待重试
        }
    }
}

方式二:Redis 记录处理状态

@Service
@Slf4j
public class OrderEventConsumer {

    private final OrderService orderService;
    private final StringRedisTemplate redisTemplate;

    @KafkaListener(topics = "order-events", groupId = "order-consumer-group")
    public void consume(ConsumerRecord<String, OrderEvent> record, 
                        Acknowledgment ack) {
        OrderEvent event = record.value();
        String dedupKey = "processed:order:" + event.getOrderId();

        // 用 SETNX 保证只处理一次
        Boolean success = redisTemplate.opsForValue()
            .setIfAbsent(dedupKey, "1", Duration.ofDays(7));

        if (Boolean.FALSE.equals(success)) {
            log.warn("消息已处理过,跳过: orderId={}", event.getOrderId());
            ack.acknowledge();
            return;
        }

        try {
            orderService.processOrderEvent(event);
            ack.acknowledge();
        } catch (Exception e) {
            // 处理失败,删除 Redis 记录,允许重试
            redisTemplate.delete(dedupKey);
            log.error("消息处理失败: orderId={}", event.getOrderId(), e);
        }
    }
}

补充说明setIfAbsent 对应 Redis 的 SETNX 命令,在 Redis 分布式环境中是原子操作,可以用来实现分布式幂等锁。

4.3 重试与死信队列(DLQ)

如果消息处理一直失败怎么办?无限重试会阻塞队列,影响其他消息的处理。

Spring Kafka 提供了 @RetryableTopic 注解,可以配置重试次数和间隔,超过重试次数后自动将消息转入死信队列(DLT)。

@Component
@Slf4j
public class OrderEventConsumer {

    private final OrderService orderService;

    /**
     * 消费订单事件
     * 重试配置:最多重试3次,每次间隔5秒
     * 3次都失败后,消息自动进入死信队列
     */
    @RetryableTopic(
        attempts = "4",  // 总尝试次数 = 1(首次)+ 3(重试)
        backoff = @Backoff(
            delay = 5000,           // 首次重试延迟5秒
            multiplier = 2.0,       // 指数退避:5s, 10s, 20s
            maxDelay = 60000        // 最大延迟60秒
        ),
        autoCreateTopics = "true",
        dltTopicSuffix = "-dlt"     // 死信队列后缀
    )
    @KafkaListener(topics = "order-events", groupId = "order-consumer-group")
    public void consume(OrderEvent event) {
        log.info("处理订单事件: orderId={}", event.getOrderId());
        orderService.processOrderEvent(event);
    }

    /**
     * 死信队列处理器
     * 所有重试失败的消息最终会进入这里
     */
    @DltHandler
    public void handleDlt(OrderEvent event) {
        log.error("订单事件进入死信队列,需要人工介入: orderId={}", 
            event.getOrderId());
        // 发送告警、记录到数据库、人工补偿
        alertService.sendAlert("订单事件处理失败,需人工处理: " + event.getOrderId());
    }
}

@RetryableTopic 的原理是:每次重试都会把消息写入一个新的重试 Topic,重试 Topic 的消费者再次尝试处理。这种非阻塞重试不会阻塞原始 Topic 的消息处理。

需要排除某些异常不触发重试(比如序列化异常,重试也没用):

@RetryableTopic(
    attempts = "4",
    backoff = @Backoff(delay = 5000),
    exclude = { SerializationException.class, 
                DeserializationException.class }  // 这些异常不重试
)
@KafkaListener(topics = "order-events", groupId = "order-consumer-group")
public void consume(OrderEvent event) {
    // ...
}

五、消费、处理、产出的事务闭环

前面提到过,Kafka 的精确一次语义有一个关键约束:它只覆盖 Kafka 到 Kafka 的路径。但如果你的场景正好是“从 Kafka 消费 → 处理 → 产出到另一个 Kafka Topic”,那 Kafka 的事务机制可以帮你做到端到端的精确一次。

5.1 配置消费者参与事务

spring:
  kafka:
    consumer:
      isolation-level: read_committed  # 只读取已提交的事务消息
    producer:
      transaction-id-prefix: tx-order-

isolation-level: read_committed 保证消费者不会读到未提交的事务消息。

5.2 消费-处理-产出的事务

@Component
@Slf4j
public class OrderEventProcessor {

    private final KafkaTemplate<String, Object> kafkaTemplate;

    @KafkaListener(topics = "order-events", groupId = "order-processor-group")
    @Transactional  // 结合 KafkaTransactionManager
    public void consume(ConsumerRecord<String, OrderEvent> record) {
        OrderEvent event = record.value();
        log.info("处理订单事件: orderId={}", event.getOrderId());

        // 1. 处理业务逻辑(比如状态转换、数据 enrichment)
        OrderProcessedEvent processedEvent = process(event);

        // 2. 发送处理结果到下游 Topic
        kafkaTemplate.send("order-processed-events", 
            processedEvent.getOrderId(), processedEvent);

        // 3. 方法正常返回,事务提交
        // 包括:Offset 提交 + 消息发送
        // 如果抛出异常,事务回滚:Offset 不提交 + 消息不发送
        log.info("订单事件处理完成,事务提交: orderId={}", event.getOrderId());
    }
}

对于 read -> process -> write 这个序列,Kafka 事务保证的是:最终对外可见的结果只有一次——要么处理结果成功写入下游 Topic 且 Offset 提交,要么什么都没发生。

但要注意:在事务回滚时,输入消息会被重新拉取,业务代码可能会被执行多次。所以,即使使用了事务,业务逻辑本身仍然需要有幂等性——因为代码可能被执行多次,只是最终只有一次提交生效。

六、性能代价:可靠性不是免费的

把可靠性拉到最高,是要付出代价的。下面这个表格可以帮助你做选型决策:

配置组合丢失风险重复风险性能影响适用场景
acks=0基准性能日志采集,允许丢失
acks=1轻微下降 (~5%)普通业务事件
acks=all + 幂等下降 ~10-15%支付、订单等关键操作
acks=all + 幂等 + 事务极低下降 ~15-20%跨Topic原子性写入

选型建议

  • 日志/监控数据acks=0acks=1,性能优先
  • 普通业务消息acks=all + 幂等,可靠性和性能的平衡点
  • 订单/支付/库存acks=all + 幂等 + 消费端幂等,双重保障
  • 跨 Topic 原子写入:在上述基础上加事务,但确认业务真的需要

七、总结

回到开头那个加了两遍积分的案例。如果当时的生产端配置了幂等性,重复消息在 Broker 层就被过滤掉了;如果消费端做了幂等设计,即使重复消息到了消费端也不会造成数据错误。

Kafka 的可靠性是一个系统工程,需要生产端和消费端共同配合:

生产端要做三件事

  1. 开启幂等性(enable.idempotence=true),防止单条消息重复
  2. 设置 acks=all,确保消息被所有副本确认
  3. 在需要原子性的场景使用事务

消费端要做三件事

  1. 关闭自动提交,手动控制 Offset
  2. 业务逻辑设计为幂等(唯一约束或状态记录)
  3. 配置重试策略和死信队列,防止无限阻塞

最后记住一句话:Kafka 的精确一次语义不覆盖数据库、不覆盖外部 API、不覆盖任何 Kafka 以外的东西。如果你的事务里调了第三方接口,自己做好补偿。

可靠性工程的本质不是“不出错”,而是“出错了能兜住”。把每一层的兜底方案都做到位了,系统才真正可靠。

系列拓展阅读

参考文献

  1. Kafka Documentation. “Exactly Once Semantics.” https://kafka.apache.org/documentation/#semantics
  2. Spring Kafka Reference Documentation. https://docs.spring.io/spring-kafka/reference/
  3. “Exactly-Once Semantics — What It Actually Means in Practice.” Trinity Logic, 2026.
  4. “Spring Boot 集成 Kafka 实战:生产者、消费者、Topic 创建与分区消费全解析.” CSDN, 2025.
  5. “Spring Boot 集成 Kafka 实战:消息生产与消费的幂等性保障.” CSDN, 2025.
  6. “Kafka 事务消息与精确一次语义理解.” CSDN, 2026.
  7. “使用 Spring @RetryableTopic 实现 Kafka 消息重试与死信队列.” 阿里云开发者社区, 2025.
  8. “Kafka 消息可靠性方案对比与实践.” CSDN, 2025.
  9. “Kafka——幂等生产者和事务生产者是一回事吗?” CSDN, 2025.
赞(0) 打赏
未经允许不得转载:MACS Dev Hub » Spring Boot + Kafka 深度实战:可靠消息传递与 Exactly-Once 语义

觉得文章有用就打赏一下文章作者

非常感谢你的打赏,我们将继续提供更多优质内容,让我们一起创建更加美好的网络世界!

支付宝扫一扫

微信扫一扫

登录

找回密码

注册