跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
📨 Kafka 3.9 教程 · 3 / 10 查看系列简介 →

Kafka 3.9 教程(二):核心概念

文章目录
  1. 目录
  2. 2.1 Kafka 的五大核心 API
  3. 2.1.1 Producer API
  4. 2.1.2 Consumer API
  5. 2.1.3 Admin API
  6. 2.1.4 Streams API
  7. 2.1.5 Connect API
  8. 2.2 生产者与消费者模型
  9. 2.2.1 生产者详解
  10. 2.2.2 消费者详解
  11. 2.2.3 生产者与消费者的解耦设计
  12. 2.3 分区机制
  13. 2.3.1 为什么需要分区?
  14. 2.3.2 分区的工作原理
  15. 2.3.3 分区数的选择
  16. 2.3.4 分区再平衡(Rebalance)
  17. 2.4 副本机制
  18. 2.4.1 副本的基本概念
  19. 2.4.2 数据同步机制
  20. 2.4.3 Leader 选举
  21. 2.4.4 副本配置最佳实践
  22. 2.4.5 副本与容错
  23. 附录:核心概念图解
  24. 生产者-消费者-主题关系
  25. 下一步学习

📌 本文是「Kafka 3.9 教程」系列第 二 篇 · 系列目录

本文档基于 Apache Kafka 3.9 官方文档整理

目录


2.1 Kafka 的五大核心 API

Kafka 提供了五套 API,涵盖了从消息发送到流处理的全流程。

2.1.1 Producer API

用于发布事件到 Kafka 主题。

主要功能:

  • 发送消息到指定主题
  • 指定消息 key 实现分区路由
  • 批量发送提高效率
  • 异步发送支持
  • 发送确认机制

基本使用:

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;

import java.util.Properties;
import java.util.concurrent.Future;

public class ProducerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        props.put("acks", "all");
        props.put("retries", 3);

        try (Producer<String, String> producer = new KafkaProducer<>(props)) {
            // 发送消息
            ProducerRecord<String, String> record =
                new ProducerRecord<>("my-topic", "key-1", "value-1");

            Future<RecordMetadata> future = producer.send(record);
            RecordMetadata metadata = future.get();

            System.out.println("消息已发送,分区:" + metadata.partition() +
                             ",偏移量:" + metadata.offset());
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

发送模式:

模式说明适用场景
fire-and-forget发送后不等待结果对可靠性要求不高的日志
同步发送等待发送结果需要确认发送成功
异步发送使用回调函数高吞吐量场景

2.1.2 Consumer API

用于订阅主题并消费事件。

主要功能:

  • 订阅一个或多个主题
  • 按分区顺序消费
  • 消费者组内负载均衡
  • 手动/自动提交偏移量
  • 支持重置偏移量

基本使用:

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class ConsumerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", "my-consumer-group");
        props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("auto.offset.reset", "earliest");
        props.put("enable.auto.commit", true);

        try (Consumer<String, String> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(Collections.singletonList("my-topic"));

            while (true) {
                ConsumerRecords<String, String> records =
                    consumer.poll(Duration.ofMillis(100));

                for (var record : records) {
                    System.out.println("分区:" + record.partition() +
                                     ",偏移量:" + record.offset() +
                                     ",键:" + record.key() +
                                     ",值:" + record.value());
                }
            }
        }
    }
}

关键配置:

配置项说明默认值
group.id消费者组 ID必填
auto.offset.reset无初始偏移量时策略latest
enable.auto.commit自动提交偏移量true
auto.commit.interval.ms自动提交间隔5000
max.poll.records每次拉取最大记录数500

2.1.3 Admin API

用于管理 Kafka 资源(主题、配置、ACL 等)。

主要功能:

  • 创建/删除主题
  • 修改主题配置
  • 查看集群元数据
  • 管理 ACL
  • 管理分区

基本使用:

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.admin.DescribeTopicsResult;
import org.apache.kafka.clients.admin.TopicDescription;

import java.util.Collections;
import java.util.Properties;
import java.util.Set;
import java.util.concurrent.ExecutionException;

public class AdminExample {
    public static void main(String[] args) throws ExecutionException, InterruptedException {
        Properties props = new Properties();
        props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");

        try (AdminClient adminClient = AdminClient.create(props)) {

            // 创建主题
            NewTopic topic = new NewTopic("my-topic", 3, (short) 1);
            adminClient.createTopics(Collections.singletonList(topic))
                      .all().get();

            System.out.println("主题创建成功");

            // 列出所有主题
            Set<String> topics = adminClient.listTopics().names().get();
            System.out.println("所有主题:" + topics);

            // 描述主题
            DescribeTopicsResult result = adminClient.describeTopics(Collections.singletonList("my-topic"));
            TopicDescription desc = result.all().get().get("my-topic");
            System.out.println("主题描述:" + desc);
        }
    }
}

2.1.4 Streams API

用于构建实时流处理应用。详见第八章:Kafka Streams。

2.1.5 Connect API

用于连接外部系统。详见第七章:Kafka Connect。


2.2 生产者与消费者模型

2.2.1 生产者详解

生产者工作流程

应用程序 → Producer → 分区选择 → 批量发送 → Broker
                        ↓
                  序列化 + 压缩

发送确认机制(acks)

acks 值含义可靠性吞吐量
0不等待确认最低最高
1等待 Leader 确认中等中等
all (-1)等待 ISR 全部确认最高最低
// 高可靠性配置
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);
props.put("min.insync.replicas", 2);
props.put("enable.idempotence", true); // 幂等性

分区策略

  1. 默认策略(轮询):均匀分布到各分区
  2. 按 key 哈希:相同 key 到相同分区
  3. 自定义分区器:实现 Partitioner 接口
// 按 key 哈希分发(保证同一用户消息到同一分区)
ProducerRecord<String, String> record =
    new ProducerRecord<>("orders", "user-123", "order data");

批量发送

// 批量大小(字节)
props.put("batch.size", 16384);  // 16KB

// 发送延迟(毫秒),0 表示立即发送
props.put("linger.ms", 5);

// 缓冲区大小
props.put("buffer.memory", 33554432);  // 32MB

压缩

// 压缩类型:none, gzip, snappy, lz4, zstd
props.put("compression.type", "lz4");

2.2.2 消费者详解

消费者组机制

核心规则:

  • 同一个消费者组内,一个分区只分配给一个消费者
  • 不同消费者组可以独立消费同一主题
  • 消费者数不应超过分区数

消费者组示意:

主题 my-topic (3 个分区)

消费者组 A:                    消费者组 B:
┌─────────────────┐           ┌─────────────────┐
│ Consumer 1 → P1 │           │ Consumer 1 → P1 │
│ Consumer 2 → P2 │           │ Consumer 2 → P2 │
│ Consumer 3 → P3 │           │ Consumer 3 → P3 │
└─────────────────┘           └─────────────────┘
(独立消费)                  (独立消费)

偏移量管理

自动提交:

props.put("enable.auto.commit", true);
props.put("auto.commit.interval.ms", 1000);

手动提交:

consumer.commitSync();  // 同步提交
consumer.commitAsync(); // 异步提交

消费策略

策略说明
earliest从最早偏移量开始消费
latest从最新偏移量开始消费
none无偏移量则报错

2.2.3 生产者与消费者的解耦设计

Kafka 的核心设计特点:生产者和消费者完全解耦。

解耦的含义:

  • 生产者不需要知道有多少消费者
  • 消费者不需要知道生产者是谁
  • 双方都不知道对方的存在
  • 可以独立扩展、升级、故障恢复

优势:

  1. 高可扩展性:可以独立扩展生产者和消费者
  2. 低耦合:系统组件间依赖最小化
  3. 高可靠性:一方故障不影响另一方
  4. 灵活架构:易于添加新的生产者和消费者

2.3 分区机制

2.3.1 为什么需要分区?

分区是 Kafka 实现水平扩展的核心机制。

需求单分区限制多分区解决方案
并行处理只能串行处理多消费者并行消费
负载均衡单点瓶颈分散到多台机器
吞吐量上限受单机限制线性扩展
存储容量受单机磁盘限制分布式存储

2.3.2 分区的工作原理

消息分配规则

  1. 指定分区:消息发送到指定分区
  2. 指定 key:按 key 哈希取模决定分区
  3. 都没有:轮询分配
// 指定分区
new ProducerRecord<>("topic", 0, "key", "value");

// 按 key 分配(保证同一用户消息有序)
new ProducerRecord<>("topic", "user-123", "order data");

分区内消息有序

保证:同一分区内的消息按发送顺序存储和消费

分区 0:
┌─────────────────────────────────────────────────┐
│ offset 0 │ offset 1 │ offset 2 │ offset 3 │ ... │
│  msg A   │  msg B   │  msg C   │  msg D   │     │
└─────────────────────────────────────────────────┘
消费顺序:A → B → C → D

跨分区无序:不同分区的消息可能交叉消费

2.3.3 分区数的选择

考虑因素

因素说明
吞吐量目标分区数影响最大并行度
消费者数量消费者数 ≤ 分区数
文件描述符限制每个分区对应多个文件
复制开销副本数 × 分区数 = 总副本数

分区数计算公式

分区数 ≈ 目标吞吐量 / 单消费者吞吐量

示例:

  • 目标吞吐量:100万消息/秒
  • 单消费者:10万消息/秒
  • 需要分区数:100万 / 10万 = 10 个分区

2.3.4 分区再平衡(Rebalance)

当消费者组成员变化时,触发分区再平衡。

触发条件:

  • 新消费者加入
  • 消费者离开(崩溃、网络断开)
  • 消费者主动退出

再平衡协议:

  1. Join:消费者请求加入组
  2. Sync:Leader 消费者分配分区,同步给组员
  3. Heartbeat:维持组关系

注意:再平衡期间,消费者暂停消费。


2.4 副本机制

2.4.1 副本的基本概念

副本类型

类型角色说明
Leader主副本处理所有读写请求
Follower从副本从 Leader 同步数据
ISR同步副本与 Leader 保持同步的副本

副本因子

# 默认副本因子
default.replication.factor=3

# 创建主题时指定
kafka-topics.sh --create --topic my-topic --replication-factor 3

2.4.2 数据同步机制

Follower 同步流程

1. Producer 发送消息到 Leader
           ↓
2. Leader 写入本地日志
           ↓
3. Follower 拉取(Fetch)数据
           ↓
4. Follower 写入本地日志
           ↓
5. Follower 发送 ACK
           ↓
6. Leader 收到 ISR 全部 ACK,提交消息

同步条件

Follower 必须满足以下条件才属于 ISR:

  1. 与 Leader 的延迟在配置时间内(replica.lag.time.max.ms)
  2. 落后消息数在阈值内(replica.lag.max.messages)

2.4.3 Leader 选举

触发时机

  • Leader 节点崩溃
  • Leader 节点主动关闭
  • 网络分区导致 Leader 不可达

选举流程

1. Controller 检测到 Leader 故障
           ↓
2. 从 ISR 中选举新 Leader
           ↓
3. 更新 ZooKeeper/KRaft 元数据
           ↓
4. 通知所有 Broker 更新
           ↓
5. 客户端重定向到新 Leader

选举原则:优先选择 ISR 中数据最新的副本

2.4.4 副本配置最佳实践

生产环境配置

# Broker 配置
default.replication.factor=3
min.insync.replicas=2

# Leader 配置
leader.imbalance.check.interval.seconds=300
leader.imbalance.per.broker.percentage=10

# Follower 配置
replica.lag.time.max.ms=10000
replica.lag.max.messages=10000

配置说明

配置项推荐值说明
replication.factor3副本因子,生产环境推荐
min.insync.replicas2最小同步副本数
replica.lag.time.max.ms10000同步超时时间
unclean.leader.election.enablefalse不允许非 ISR 选举

2.4.5 副本与容错

不同副本因子的容错能力

副本因子可容忍故障数说明
10无容错,单点故障
20允许 1 副本故障,但无法选举
31允许 1 个节点故障

故障场景处理

场景1:1个 Follower 故障

  • 自动从 ISR 移除
  • 不影响 Leader 工作
  • 故障恢复后追赶同步

场景2:Leader 故障

  • 从 ISR 选举新 Leader
  • 旧 Leader 恢复后成为 Follower
  • 可能有短暂不可用(毫秒级)

场景3:多个节点故障

  • 若故障数 ≥ ISR 数量,无法选举
  • 集群不可用,等待节点恢复

附录:核心概念图解

生产者-消费者-主题关系

                    ┌─────────────────────────────────────────┐
                    │              Kafka 集群                  │
                    │                                         │
┌──────────────┐    │  ┌─────────────────────────────────┐   │
│   生产者 1   │───▶│  │  主题: orders                   │   │
├──────────────┤    │  │  ┌─────┬─────┬─────┬─────┐     │   │
│   生产者 2   │───▶│  │  │ P0  │ P1  │ P2  │ P3  │     │   │
├──────────────┤    │  │  └─────┴─────┴─────┴─────┘     │   │
│   生产者 N   │───▶│  │    │     │     │     │         │   │
└──────────────┘    │  │    ▼     ▼     ▼     ▼         │   │
                    │  │  ┌─────┬─────┬─────┬─────┐     │   │
                    │  │  │Broker│Broker│Broker│Broker│     │   │
                    │  │  │  #1  │  #2  │  #3  │  #4  │     │   │
                    │  │  └─────┴─────┴─────┴─────┘     │   │
                    │  └─────────────────────────────────┘   │
                    └─────────────────────────────────────────┘
                                        │
                    ┌───────────────────┴───────────────────┐
                    ▼                                       ▼
┌──────────────────────────────┐     ┌──────────────────────────────┐
│        消费者组 A             │     │        消费者组 B             │
│  ┌─────┐ ┌─────┐ ┌─────┐    │     │  ┌─────┐ ┌─────┐ ┌─────┐    │
│  │ C1  │ │ C2  │ │ C3  │    │     │  │ C1  │ │ C2  │ │ C3  │    │
│  └─────┘ └─────┘ └─────┘    │     │  └─────┘ └─────┘ └─────┘    │
│      ↓      ↓      ↓        │     │      ↓      ↓      ↓        │
│      P0     P1     P2       │     │      P0     P1     P2       │
└──────────────────────────────┘     └──────────────────────────────┘
(每组独立消费)                  (每组独立消费)

下一步学习


RAG 智能问答

针对本文继续提问:《Kafka 3.9 教程(二):核心概念》