Kafka可靠性幂等与事务
Kafka 可靠性不能只回答“Kafka 会不会丢消息”。正确的拆法是看三段链路:生产者有没有可靠发送,Broker 有没有可靠存储,消费者有没有可靠处理。
flowchart TD
A["生产阶段"] --> B["消息是否成功写入 Kafka"]
C["存储阶段"] --> D["Broker 宕机是否丢数据"]
E["消费阶段"] --> F["业务处理和 Offset 是否一致"]任何一段没设计好,都可能造成丢消息、重复消息、乱序或数据不一致。
生产阶段可靠性
生产者发送消息时,常见关键配置是:
| 配置 | 建议 | 作用 |
|---|---|---|
acks | 重要业务用 all | 等待 Leader 和同步副本确认 |
retries | 开启重试 | 网络抖动时自动重发 |
enable.idempotence | 开启 | 避免生产者重试导致分区内重复写 |
delivery.timeout.ms | 合理设置 | 控制发送最终超时时间 |
linger.ms | 适当增大 | 等待更多消息组成批次,提高吞吐 |
acks 的含义
| acks | 含义 | 风险 |
|---|---|---|
0 | 生产者不等 Broker 确认 | 吞吐高,但消息可能直接丢 |
1 | Leader 写入成功就返回 | Leader 宕机且 Follower 未同步时可能丢 |
all | 等待 ISR 中足够副本确认 | 更可靠,但延迟更高 |
flowchart TD
A["Producer 发送"] --> B["Leader 写入"]
B --> C{"acks 配置"}
C -->|"0"| D["不等待确认"]
C -->|"1"| E["Leader 写入即成功"]
C -->|"all"| F["等待同步副本确认"]生产环境里,核心业务通常使用 acks=all 和 enable.idempotence=true。但只配这两个还不够,还要看 Broker 的最小同步副本。
存储阶段可靠性
Kafka 的副本机制依赖 Leader、Follower 和 ISR。生产者写 Leader,Follower 从 Leader 复制。
flowchart TD
A["Producer"] --> B["Leader"]
B --> C["Follower 1"]
B --> D["Follower 2"]
C --> E["ISR"]
D --> E
E --> F["满足确认条件后返回成功"]关键配置:
| 配置 | 作用 |
|---|---|
replication.factor | 每个分区有几个副本 |
min.insync.replicas | 至少多少同步副本确认才算成功 |
unclean.leader.election.enable | 是否允许落后副本当 Leader |
例如 3 副本场景:
replication.factor=3min.insync.replicas=2- 生产者
acks=all
这表示至少 Leader 和一个同步副本确认后,生产者才认为成功。这样允许一个 Broker 宕机时仍尽量不丢已确认消息。
如果 min.insync.replicas=1,即使 acks=all,实际也可能只有 Leader 自己确认。一旦 Leader 还没复制就宕机,仍然有丢失风险。
消费阶段可靠性
消费端最常见的问题是 offset 提交时机。
错误流程
flowchart TD
A["Consumer 拉取消息"] --> B["自动提交 Offset"]
B --> C["处理数据库"]
C --> D["服务宕机"]
D --> E["重启后从新 Offset 消费"]
E --> F["旧消息漏处理"]推荐流程
flowchart TD
A["Consumer 拉取消息"] --> B["执行业务逻辑"]
B --> C{"业务成功?"}
C -->|"是"| D["提交 Offset"]
C -->|"否"| E["抛异常或进入重试"]业务成功后再提交 offset,可以避免漏消费,但会带来重复消费可能。因此消费端幂等是必须的。
消费幂等
幂等的意思是:同一条消息执行一次和执行多次,最终业务结果一致。
常见方案:
| 方案 | 适合场景 | 原理 |
|---|---|---|
| 数据库唯一索引 | 订单、支付、流水 | 同一个业务 key 只能插入一次 |
| 状态机判断 | 订单状态流转 | 只有合法状态才能推进 |
Redis setnx | 短期去重 | 第一次处理成功写入 key |
| 消费日志表 | 通用消费记录 | groupId + messageKey 唯一 |
示例:用数据库唯一索引做消费日志。
create table kafka_consume_log (
id bigint primary key auto_increment,
consumer_group varchar(128) not null,
topic varchar(128) not null,
message_key varchar(128) not null,
created_at datetime not null,
unique key uk_consume (consumer_group, topic, message_key)
);@Transactional
public void handle(OrderCreatedEvent event) {
boolean first = consumeLogRepository.tryInsert(
"search-index-service",
"order.created",
String.valueOf(event.orderId())
);
if (!first) {
return;
}
searchIndexService.rebuild(event.orderId());
}重试和死信
业务异常不能无限重试,否则会阻塞分区后续消息。常见设计:
- 可恢复异常:重试几次,例如数据库短暂不可用。
- 不可恢复异常:直接进入死信 Topic,例如参数格式错误。
- 告警和人工处理:死信消息需要监控,不能只是换个 Topic 放着。
Spring Kafka 中可以配置死信:
@Configuration
public class KafkaErrorHandlerConfig {
@Bean
public DefaultErrorHandler defaultErrorHandler(KafkaTemplate<Object, Object> template) {
DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(
template,
(record, ex) -> new TopicPartition(record.topic() + ".DLT", record.partition())
);
return new DefaultErrorHandler(recoverer, new FixedBackOff(1000L, 3L));
}
}这个配置表示:消费失败后间隔 1 秒重试,最多重试 3 次,仍失败就发送到原 Topic 对应的 .DLT 死信 Topic。
事务和 Exactly Once
Kafka 支持生产者幂等和事务,但要区分几个层次:
| 层次 | 解决什么 | 不能解决什么 |
|---|---|---|
| 生产者幂等 | 同一 Producer 重试不重复写入同一分区 | 业务重复调用仍可能重复发 |
| Kafka 事务 | 多个分区写入、消费后再生产的原子性 | 不能自动保证外部数据库事务一致 |
| 消费端幂等 | 重复投递时业务结果不重复 | 需要业务自己设计唯一键或状态机 |
Kafka 的 Exactly Once 更准确地说,是 Kafka 流处理链路内的精确一次语义。比如“消费 Topic A,处理后写 Topic B”,可以用事务保证读取位置和输出结果一致。
但如果业务同时写 MySQL 和 Kafka,就不是 Kafka 单独能解决的。常用方案是:
- 本地消息表 Outbox。
- CDC 采集数据库变更写 Kafka。
- 分布式事务框架。
- 业务补偿和对账。
Outbox 可靠发送模式
Outbox 是商业项目里非常常用的方案:业务数据和消息记录在同一个本地事务内落库,然后由后台任务或 CDC 把消息发送到 Kafka。
flowchart TD
A["业务请求"] --> B["本地事务"]
B --> C["写业务表"]
B --> D["写 outbox_message 表"]
D --> E["后台任务或 CDC 扫描"]
E --> F["发送 Kafka"]
F --> G["标记消息已发送"]这样可以避免“数据库提交成功,但 Kafka 发送失败”的不一致。
示例表:
create table outbox_message (
id bigint primary key auto_increment,
aggregate_id varchar(128) not null,
topic varchar(128) not null,
message_key varchar(128) not null,
payload json not null,
status varchar(32) not null,
created_at datetime not null,
updated_at datetime not null
);可靠性检查清单
| 阶段 | 必查项 |
|---|---|
| 生产者 | acks=all、重试、幂等、发送失败告警 |
| Broker | 副本数、最小同步副本、磁盘容量、ISR 监控 |
| Topic | 分区数、保留时间、压缩策略、死信 Topic |
| 消费者 | 手动提交、幂等、重试、死信、消费延迟监控 |
| 业务 | 唯一键、状态机、补偿、对账 |
小结
Kafka 可靠性不是一个参数,而是一套链路设计。生产者保证可靠写入,Broker 保证副本存储,消费者保证业务成功后提交 offset,业务系统保证幂等和补偿。只要其中一环偷懒,线上就会用重复、丢失、乱序来提醒你。
