Skip to content
第 13 章 ⏱ 14 分钟阅读

第 13 章:Kafka 进阶 ​

学习目标 ​

  • 理解 Kafka 的分区、副本、ISR 机制
  • 掌握生产者/消费者 Java 代码
  • 处理消息顺序、幂等、事务
  • 监控 Kafka 集群状态

一、Kafka 核心概念 ​

概念含义
Producer消息生产者
Consumer消息消费者
Broker一个 Kafka 进程
Topic消息主题
Partition分区,物理存储单位
Offset消息在分区内的编号
Consumer Group消费者组,组内均分分区
ISRIn-Sync Replicas,与 Leader 同步的副本集合

二、为什么 Kafka 这么快 ​

  • 顺序写盘:append-only,机械盘也能百万 QPS
  • 零拷贝:sendfile 系统调用,跳过用户态
  • 页缓存:不落盘,直接用 OS 缓存
  • 批量发送:攒一批再发,降低网络开销

三、安装 ​

bash
docker run -d --name kafka \
  -p 9092:9092 \
  -e KAFKA_BROKER_ID=1 \
  -e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
  -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
  apache/kafka:3.7.0

(新版可用 KRaft 模式,不依赖 ZooKeeper)

四、核心原理:分区与顺序 ​

同一个分区内的消息有序,不同分区不保证全局顺序。

text
Topic: orders (3 partitions)
┌─────────────┬─────────────┬─────────────┐
│ Partition 0 │ Partition 1 │ Partition 2 │
│  offset 0   │  offset 0   │  offset 0   │
│  offset 1   │  offset 1   │  offset 1   │
└─────────────┴─────────────┴─────────────┘

生产时指定 partition key → 同 key 的消息路由到同一分区 → 保证同用户/同订单有序。

⚠️ 坑 1:全局顺序 = topic 只设 1 个分区 = 失去并发能力。生产一般分区内有序就够。

五、生产者(Java) ​

xml
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>
java
@Resource
private KafkaTemplate<String, OrderEvent> kafkaTemplate;

public void send(OrderEvent event) {
    CompletableFuture<SendResult<String, OrderEvent>> future =
        kafkaTemplate.send("orders", String.valueOf(event.getUserId()), event);

    future.whenComplete((result, ex) -> {
        if (ex != null) {
            log.error("发送失败", ex);
            // 补偿:落库 + 定时重试
        }
    });
}
yaml
spring:
  kafka:
    producer:
      acks: all                    # 等待所有 ISR 确认
      retries: 3
      batch-size: 16384
      properties:
        enable.idempotence: true   # ⚠️ 坑 2:生产必须开幂等

⚠️ 坑 2:enable.idempotence=true 保证不丢不重(生产者侧);Kafka 0.11+ 默认开启,别手动关。

六、消费者(Java) ​

java
@Component
public class OrderConsumer {

    @KafkaListener(topics = "orders", groupId = "order-service")
    public void handle(ConsumerRecord<String, OrderEvent> record,
                       @Header(KafkaHeaders.OFFSET) long offset) {
        try {
            // 业务处理
            log.info("收到订单:{}", record.value());
        } catch (Exception e) {
            // 抛出 → Spring Kafka 自动重试 → 失败进 DLT
            throw new RuntimeException(e);
        }
    }
}
yaml
spring:
  kafka:
    consumer:
      group-id: order-service
      auto-offset-reset: earliest      # ⚠️ 坑 3:新组从最早开始,避免漏消息
      enable-auto-commit: false       # 手动提交
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer

    listener:
      ack-mode: manual_immediate
      concurrency: 3                   # 并发消费数 ≤ 分区数

⚠️ 坑 3:auto-offset-reset=latest 新组会丢历史消息;想补数据用 earliest。 ⚠️ 坑 4:concurrency 不能超过分区数,否则多出来的线程空闲。

七、消费者组与 Rebalance ​

  • 一个 partition 同一时刻只被组内一个消费者消费
  • 消费者数 > 分区数 → 多余消费者空闲
  • 消费者宕机 → 触发 rebalance,分区重新分配

⚠️ 坑 5:MAX_POLL_INTERVAL_MS 默认 5 分钟 —— 处理慢的任务要调大,否则被踢出组,反复 rebalance。

八、幂等消费 ​

java
public void handle(OrderEvent event) {
    String key = "order:" + event.getId();
    if (redis.opsForValue().setIfAbsent(key, "1", Duration.ofHours(24))) {
        // 第一次消费,处理
        process(event);
    } else {
        log.warn("重复消息:{}", event.getId());
    }
}

九、Kafka 事务(精确一次) ​

java
@Bean
public ProducerFactory<String, OrderEvent> producerFactory() {
    Map<String, Object> props = new HashMap<>();
    props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-tx-1");
    // ...
    return new DefaultKafkaProducerFactory<>(props);
}
java
kafkaTemplate.executeInTransaction(ops -> {
    ops.send("orders", event);
    ops.send("audit", event);    // 两个 topic 原子提交
    return null;
});

⚠️ 坑 6:TRANSACTIONAL_ID 必须全局唯一,不然两个生产者用同一 ID 会相互干扰。

十、监控 ​

bash
# Kafka 内置工具
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group order-service --describe

# 看消费进度、lag

Lag 持续增大 → 消费慢,要扩容消费者或优化业务。

十一、本章小结 ​

要点关键
分区topic 物理切分,同分区有序
acks=all + 幂等生产侧不丢不重
消费者组同组内均分分区
offsetearliest 防止新组丢消息
幂等Redis SETNX 判重
事务跨 topic / 跨分区原子
lag监控消费延迟,及时扩容

动手练习 ​

  1. 用 Docker 启动 Kafka(单 broker + 3 分区 topic)
  2. 生产者用 userId 当 key,确保同一用户的消息落到同分区
  3. 启动两个消费者(同一 group),观察分区被均分

下一章:第 14 章:RocketMQ →

本站基于 VitePress 构建 · 由 StackHub 团队维护