从消息丢失到精确一次处理,一篇讲透生产级 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 | 不等待任何确认 | 最高性能,消息可能丢 |
1 | Leader 副本写入即确认 | 折中方案,Leader 宕机可能丢 |
all(或 -1) | ISR 中所有副本都写入才确认 | 最高可靠性 |
retries:发送失败后的重试次数。生产环境建议设成一个很大的值(如 Integer.MAX_VALUE),让重试机制配合幂等性来保证最终成功。
enable.idempotence:幂等生产者开关。开启后,Broker 会给每个生产者分配一个 PID(Producer ID),并为每个分区维护一个序列号。如果生产者重试一条已经写成功的消息,Broker 看到相同的序列号就直接丢弃,不会重复写入。
注意:开启幂等性后,
acks会被自动设为all,retries会被设为Integer.MAX_VALUE,max.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 发送消息时的异常处理
KafkaTemplate 的 send() 方法是异步的,返回 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=0或acks=1,性能优先 - 普通业务消息:
acks=all+ 幂等,可靠性和性能的平衡点 - 订单/支付/库存:
acks=all+ 幂等 + 消费端幂等,双重保障 - 跨 Topic 原子写入:在上述基础上加事务,但确认业务真的需要
七、总结
回到开头那个加了两遍积分的案例。如果当时的生产端配置了幂等性,重复消息在 Broker 层就被过滤掉了;如果消费端做了幂等设计,即使重复消息到了消费端也不会造成数据错误。
Kafka 的可靠性是一个系统工程,需要生产端和消费端共同配合:
生产端要做三件事:
- 开启幂等性(
enable.idempotence=true),防止单条消息重复 - 设置
acks=all,确保消息被所有副本确认 - 在需要原子性的场景使用事务
消费端要做三件事:
- 关闭自动提交,手动控制 Offset
- 业务逻辑设计为幂等(唯一约束或状态记录)
- 配置重试策略和死信队列,防止无限阻塞
最后记住一句话:Kafka 的精确一次语义不覆盖数据库、不覆盖外部 API、不覆盖任何 Kafka 以外的东西。如果你的事务里调了第三方接口,自己做好补偿。
可靠性工程的本质不是“不出错”,而是“出错了能兜住”。把每一层的兜底方案都做到位了,系统才真正可靠。
系列拓展阅读
- 《Redis 分布式缓存实战:从 Spring Cache 到缓存雪崩/穿透/击穿解决方案》 —— 消息队列与缓存的协同
- 《JUnit 5 + Testcontainers:Java 微服务集成测试最佳实践》 —— 用 Testcontainers 做 Kafka 集成测试
- 《Java 应用全链路追踪实战:OpenTelemetry + Jaeger 从入门到生产》 —— 消息链路中的追踪实践
- 《Spring Boot 微服务间调用最佳实践:从 OpenFeign 到 HttpExchange + 服务治理》 —— 同步调用与异步消息的对比
参考文献
- Kafka Documentation. “Exactly Once Semantics.” https://kafka.apache.org/documentation/#semantics
- Spring Kafka Reference Documentation. https://docs.spring.io/spring-kafka/reference/
- “Exactly-Once Semantics — What It Actually Means in Practice.” Trinity Logic, 2026.
- “Spring Boot 集成 Kafka 实战:生产者、消费者、Topic 创建与分区消费全解析.” CSDN, 2025.
- “Spring Boot 集成 Kafka 实战:消息生产与消费的幂等性保障.” CSDN, 2025.
- “Kafka 事务消息与精确一次语义理解.” CSDN, 2026.
- “使用 Spring @RetryableTopic 实现 Kafka 消息重试与死信队列.” 阿里云开发者社区, 2025.
- “Kafka 消息可靠性方案对比与实践.” CSDN, 2025.
- “Kafka——幂等生产者和事务生产者是一回事吗?” CSDN, 2025.









