java 技术随笔

Kafka 核心原理与 Spring Boot 实战:分区、副本、消费组与 offset 提交

Kafka 是当下最主流的高吞吐消息中间件:日志采集用它、实时计算用它、订单削峰也可能用它。本文从 Topic/Partition/Replica/Offset 这几个核心概念讲起,说清楚一条消息从生产到消费的完整链路,再给出 Spring Boot 接入的完整示例,最后总结消费组、offset 提交等最容易踩的坑。

核心概念:一张图看懂 Kafka 数据组织

Topic(主题):一类消息的逻辑名字,比如 order_topic。一条消息一定属于某个 Topic。

Partition(分区):Topic 被物理切成多个分区,每个分区是一个有序的日志文件,追加写入。分区是 Kafka 并行和水平扩展的根本:分区越多,并行度越高。

Replica(副本):每个分区有多个副本,分 Leader 和 Follower。读写都走 Leader,Follower 异步同步数据,Leader 挂了从 ISR(同步副本集合)里选新 Leader,保证不丢数据。

Offset(偏移量):分区内每条消息的位置编号,从 0 递增。消费组消费到哪了,就是记 offset 到哪。

Consumer Group(消费组):一组消费者共同消费一个 Topic。一条消息在一个组内只会被一个消费者实例消费;不同组互不影响,各记各的 offset。这就是"广播给多个系统、组内负载均衡"的原理。

早期版本 Kafka 的元数据(分区、副本、Leader 信息)靠 ZooKeeper 保存,Kafka 3.x 之后用自研的 KRaft 模式(内部 Raft 协议)管理,可以完全去掉 ZooKeeper,部署简单一大截。

生产链路:消息是怎么写进去的

Producer 发消息默认流程:

1. 找分区:消息按 key 的 hash 取模选分区,同一 key 进同一分区(保证有序);没 key 则按粘性分区策略轮询,尽量攒批。可以显式指定 partition,或自定义 Partitioner。

2. 攒批发送:Producer 端有缓冲区,按 batch.size 和 linger.ms 攒一批再发,配合压缩(compression.type=gzip/lz4/zstd),把网络和 IO 效率拉满,这是 Kafka 吞吐高的关键之一。

3. Broker 落盘:消息追加写到分区日志(顺序 IO,非常快),并维护到 ISR。Follower 从 Leader 拉取同步。

4. 确认回执:由 acks 决定什么时候算发送成功——acks=0 发出去不管;acks=1 Leader 写成功就返回(可能丢:Leader 挂了没同步);acks=all 要等 ISR 里所有副本都写成功才返回(最可靠,配合 min.insync.replicas=2 保证至少两个副本)。

消费端是拉模式(pull),消费者主动向 Broker 拉数据,Broker 不推,所以 Kafka 天然能扛住消费慢的情况,代价是会有轻微延迟。

消费链路:消费组与 offset 提交

一个消费组内多个消费者会做分区分配(Range/RoundRobin/Sticky 策略),新成员加入或挂掉都会触发重平衡(Rebalance),期间该组暂停消费。常见坑:消费线程数大于分区数时,多的线程闲着没数据;消费很慢触发超时被踢出组,又触发新一轮 Rebalance,形成"踢出-再平衡"恶性循环。

offset 提交分两种:自动提交(enable.auto.commit=true,默认 5 秒一次)在 poll 时异步提交,可能重复消费;手动提交(enable.auto.commit=false)在业务处理完成后同步/异步提交,可以做到消息"处理完才提交",配合业务幂等实现至少一次投递。生产环境建议手动提交。

消费的关键注意点:先处理业务再提交 offset,否则消费到一半进程重启,offset 已提交会丢消息;反过来先提交后处理会重复消费,需幂等兜底。

Spring Boot 集成示例

引入 spring-kafka 依赖后,生产者用 KafkaTemplate 发送:

// 生产者:发送消息,指定 key 保证同一订单消息有序
@Service
public class OrderProducer {
    private final KafkaTemplate<String, Object> kafkaTemplate;

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

    public void sendOrderCreated(Long orderId) {
        OrderEvent event = new OrderEvent(orderId, "CREATED");
        // 同一 key 进同一分区,保证该订单消息顺序消费
        kafkaTemplate.send("order-topic", String.valueOf(orderId), event);
    }
}

消费者用 @KafkaListener 注解即可,配合手动 ack:

// 消费者:手动提交 offset,业务成功后再 ack
@Component
public class OrderConsumer {

    @KafkaListener(topics = "order-topic", groupId = "order-group")
    public void onOrder(ConsumerRecord<String, Object> record,
                        Acknowledgment ack) {
        try {
            OrderEvent event = (OrderEvent) record.value();
            // 业务幂等:先去重表查一次,处理过直接跳过
            if (orderService.isProcessed(event.getOrderId())) {
                ack.acknowledge();
                return;
            }
            orderService.handleOrderCreated(event);
            // 业务成功,标记幂等 + 提交 offset
            orderService.markProcessed(event.getOrderId());
            ack.acknowledge();
        } catch (Exception e) {
            // 记录日志/进死信,不 ack 等待重试或人工处理
            log.error("消费失败", e);
        }
    }
}

application.yml 关键配置:

spring:
  kafka:
    bootstrap-servers: 192.168.1.10:9092
    producer:
      acks: all          # 所有副本确认,防丢
      retries: 3
      compression-type: gzip
    consumer:
      group-id: order-group
      enable-auto-commit: false   # 手动提交
      auto-offset-reset: earliest # 无 offset 时从最早开始
    listener:
      ack-mode: manual    # 手动 ack

高频坑与面试速答

Q:怎么保证消息不丢?生产端 acks=all + min.insync.replicas=2,消费端手动提交、先业务后 ack,Broker 端关闭 unclean.leader.election(不让落后太多的副本当 Leader)。

Q:怎么保证有序?同一业务 key 发到同一分区 + 单分区内严格有序 + 该分区只由一个消费者线程消费。

Q:重复消费怎么办?Kafka 至少一次语义下重复必然存在,消费端做幂等(唯一键、状态机、去重表)。

Q:消费积压了怎么处理?先看是"消费不过来"还是"卡死重平衡":前者扩容消费者(不超过分区数)或加分区,后者查日志找异常消息,必要时跳过问题消息。

一句话记住 Kafka:分区是并行单元,副本是可靠性保障,offset 是消费进度,消费组是广播与负载均衡的载体。

标签
Kafka消息队列Spring Boot高并发