Skip to content
第 198 / 250 章中间件⏱ 14 分钟阅读

第 198 章:Kafka 入门与架构

学习目标

  • 理解 Kafka 核心概念
  • 掌握分区、副本、消费者组
  • 学会基础 Producer / Consumer
  • 了解 Kafka 高可用机制

一、Kafka 是什么

Kafka 是一个分布式流处理平台,核心能力:

  • 消息队列:发布订阅
  • 持久化存储:磁盘存储,可重放
  • 流处理:Kafka Streams

二、Kafka 核心概念

2.1 Broker

Kafka 集群中的一台服务器就是一个 broker。

2.2 Topic(主题)

消息分类,生产者发送到 topic,消费者从 topic 读取。

2.3 Partition(分区)

  • 一个 topic 可分多个 partition
  • 分区内消息有序,跨分区不保证
  • 每个分区是一个有序的追加日志

2.4 Offset(偏移量)

分区内的消息编号,消费者通过 offset 标记消费位置。

Partition 0:
[0] Hello    ← offset 0
[1] World    ← offset 1
[2] Kafka    ← offset 2
[3] ...      ← offset 3

2.5 Replica(副本)

每个分区有多个副本:

  • Leader:处理读写
  • Follower:同步数据,Leader 挂了顶上

2.6 Consumer Group(消费者组)

  • 同一 group 内,一个 partition 只能被一个消费者消费
  • 多个 partition 在 group 内负载均衡
  • 不同 group 独立消费同一份数据(广播)

三、Kafka 架构

四、安装

4.1 Docker Compose

yaml
# docker-compose.yml
version: '3.8'
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181

  kafka:
    image: confluentinc/cp-kafka:7.5.0
    depends_on: [zookeeper]
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
bash
docker-compose up -d

4.2 CLI 命令

bash
# 进容器
docker exec -it kafka bash

# 创建 topic
kafka-topics --bootstrap-server localhost:9092 \
  --create --topic test --partitions 3 --replication-factor 1

# 列出
kafka-topics --bootstrap-server localhost:9092 --list

# 详情
kafka-topics --bootstrap-server localhost:9092 --describe --topic test

# 发送
kafka-console-producer --bootstrap-server localhost:9092 --topic test
> hello
> world

# 消费
kafka-console-consumer --bootstrap-server localhost:9092 --topic test --from-beginning

五、Spring Boot Kafka

5.1 引入

xml
<dependency>
  <groupId>org.springframework.kafka</groupId>
  <artifactId>spring-kafka</artifactId>
</dependency>

5.2 配置

yaml
spring:
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      acks: all
      retries: 3
      properties:
        enable.idempotence: true
    consumer:
      group-id: my-group
      auto-offset-reset: earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.example.dto"

5.3 Producer

java
@Service
public class OrderProducer {

    @Autowired
    private KafkaTemplate<String, OrderEvent> kafkaTemplate;

    public void send(OrderEvent event) {
        // 用订单 ID 作 key,保证同一订单到同一分区
        kafkaTemplate.send("order-events", event.getOrderId().toString(), event);
    }

    // 带回调
    public void sendWithCallback(OrderEvent event) {
        CompletableFuture<SendResult<String, OrderEvent>> future =
            kafkaTemplate.send("order-events", event.getOrderId().toString(), event);

        future.whenComplete((result, ex) -> {
            if (ex == null) {
                log.info("发送成功: {}", result.getRecordMetadata().offset());
            } else {
                log.error("发送失败", ex);
            }
        });
    }
}

5.4 Consumer

java
@Component
public class OrderConsumer {

    @KafkaListener(topics = "order-events", groupId = "inventory-group")
    public void handleOrder(OrderEvent event,
                            @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                            @Header(KafkaHeaders.OFFSET) long offset) {
        log.info("收到订单事件: partition={}, offset={}, event={}",
                 partition, offset, event);

        // 业务处理
        inventoryService.deduct(event);
    }

    // 批量消费
    @KafkaListener(topics = "batch-events", groupId = "batch-group")
    public void handleBatch(List<OrderEvent> events) {
        inventoryService.batchDeduct(events);
    }
}

六、消息可靠性

6.1 Producer 可靠性

yaml
spring:
  kafka:
    producer:
      acks: all                  # 所有 ISR 确认
      retries: 3                 # 重试次数
      properties:
        enable.idempotence: true # 幂等性,防重复
        max.in.flight.requests.per.connection: 5
acks含义可靠性
0不等确认可能丢
1Leader 写入即返回极少数丢
all所有 ISR 写入才返回不丢

6.2 Consumer 可靠性

java
// 手动提交 offset
@KafkaListener(topics = "order-events", groupId = "inventory-group")
public void handleOrder(OrderEvent event, Acknowledgment ack) {
    try {
        inventoryService.deduct(event);
        ack.acknowledge();   // 提交 offset
    } catch (Exception e) {
        // 不提交,下次重试
    }
}
yaml
spring:
  kafka:
    consumer:
      enable-auto-commit: false  # 关闭自动提交
    listener:
      ack-mode: manual_immediate

6.3 至少一次(最常用)

业务处理成功 + 提交 offset = 至少一次

重复消费

可能因失败重试导致重复消费,业务方需幂等

6.4 恰好一次(很难)

依赖 Kafka 事务 + 幂等 Producer + 消费事务,复杂且性能低,实际少用。

七、消息顺序性

7.1 分区内有序

Kafka 只保证分区内有序,跨分区不保证。

7.2 同一订单顺序

java
// 用订单 ID 作 key,Hash 到同一分区
kafkaTemplate.send("order-events", orderId.toString(), event);

7.3 全局有序(不推荐)

bash
kafka-topics --create --topic xxx --partitions 1

只用一个分区,但牺牲并发。

八、Consumer Group 重平衡

消费者组中,消费者加入或退出会触发rebalance(重新分配分区)。

避免频繁 rebalance

  • 合理设置 session.timeout.ms
  • 实现 ConsumerRebalanceListener 保存状态

九、本章小结

概念作用
Broker服务器节点
Topic消息分类
Partition分区(并行单位)
Offset消息位置
Consumer Group消费者组
ISR同步副本集合
可靠性acks重复性能
最多一次0最快
至少一次all
恰好一次all + 事务

动手练习

  1. 用 Docker Compose 启动 Kafka
  2. 用 Spring Boot 实现 Producer / Consumer
  3. 手动提交 offset 实现至少一次
  4. 用 key 保证同一订单顺序

推荐阅读


下一章:第 199 章:Kafka 进阶与可靠性

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