Skip to content

Kafka可靠性幂等与事务

Kafka 可靠性不能只回答“Kafka 会不会丢消息”。正确的拆法是看三段链路:生产者有没有可靠发送,Broker 有没有可靠存储,消费者有没有可靠处理。

mermaid
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 确认吞吐高,但消息可能直接丢
1Leader 写入成功就返回Leader 宕机且 Follower 未同步时可能丢
all等待 ISR 中足够副本确认更可靠,但延迟更高
mermaid
flowchart TD
    A["Producer 发送"] --> B["Leader 写入"]
    B --> C{"acks 配置"}
    C -->|"0"| D["不等待确认"]
    C -->|"1"| E["Leader 写入即成功"]
    C -->|"all"| F["等待同步副本确认"]

生产环境里,核心业务通常使用 acks=allenable.idempotence=true。但只配这两个还不够,还要看 Broker 的最小同步副本。

存储阶段可靠性

Kafka 的副本机制依赖 Leader、Follower 和 ISR。生产者写 Leader,Follower 从 Leader 复制。

mermaid
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=3
  • min.insync.replicas=2
  • 生产者 acks=all

这表示至少 Leader 和一个同步副本确认后,生产者才认为成功。这样允许一个 Broker 宕机时仍尽量不丢已确认消息。

如果 min.insync.replicas=1,即使 acks=all,实际也可能只有 Leader 自己确认。一旦 Leader 还没复制就宕机,仍然有丢失风险。

消费阶段可靠性

消费端最常见的问题是 offset 提交时机。

错误流程

mermaid
flowchart TD
    A["Consumer 拉取消息"] --> B["自动提交 Offset"]
    B --> C["处理数据库"]
    C --> D["服务宕机"]
    D --> E["重启后从新 Offset 消费"]
    E --> F["旧消息漏处理"]

推荐流程

mermaid
flowchart TD
    A["Consumer 拉取消息"] --> B["执行业务逻辑"]
    B --> C{"业务成功?"}
    C -->|"是"| D["提交 Offset"]
    C -->|"否"| E["抛异常或进入重试"]

业务成功后再提交 offset,可以避免漏消费,但会带来重复消费可能。因此消费端幂等是必须的。

消费幂等

幂等的意思是:同一条消息执行一次和执行多次,最终业务结果一致。

常见方案:

方案适合场景原理
数据库唯一索引订单、支付、流水同一个业务 key 只能插入一次
状态机判断订单状态流转只有合法状态才能推进
Redis setnx短期去重第一次处理成功写入 key
消费日志表通用消费记录groupId + messageKey 唯一

示例:用数据库唯一索引做消费日志。

sql
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)
);
java
@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 中可以配置死信:

java
@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。

mermaid
flowchart TD
    A["业务请求"] --> B["本地事务"]
    B --> C["写业务表"]
    B --> D["写 outbox_message 表"]
    D --> E["后台任务或 CDC 扫描"]
    E --> F["发送 Kafka"]
    F --> G["标记消息已发送"]

这样可以避免“数据库提交成功,但 Kafka 发送失败”的不一致。

示例表:

sql
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,业务系统保证幂等和补偿。只要其中一环偷懒,线上就会用重复、丢失、乱序来提醒你。