📌 本文是「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); // 幂等性
分区策略
- 默认策略(轮询):均匀分布到各分区
- 按 key 哈希:相同 key 到相同分区
- 自定义分区器:实现 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 的核心设计特点:生产者和消费者完全解耦。
解耦的含义:
- 生产者不需要知道有多少消费者
- 消费者不需要知道生产者是谁
- 双方都不知道对方的存在
- 可以独立扩展、升级、故障恢复
优势:
- 高可扩展性:可以独立扩展生产者和消费者
- 低耦合:系统组件间依赖最小化
- 高可靠性:一方故障不影响另一方
- 灵活架构:易于添加新的生产者和消费者
2.3 分区机制
2.3.1 为什么需要分区?
分区是 Kafka 实现水平扩展的核心机制。
| 需求 | 单分区限制 | 多分区解决方案 |
|---|---|---|
| 并行处理 | 只能串行处理 | 多消费者并行消费 |
| 负载均衡 | 单点瓶颈 | 分散到多台机器 |
| 吞吐量上限 | 受单机限制 | 线性扩展 |
| 存储容量 | 受单机磁盘限制 | 分布式存储 |
2.3.2 分区的工作原理
消息分配规则
- 指定分区:消息发送到指定分区
- 指定 key:按 key 哈希取模决定分区
- 都没有:轮询分配
// 指定分区
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)
当消费者组成员变化时,触发分区再平衡。
触发条件:
- 新消费者加入
- 消费者离开(崩溃、网络断开)
- 消费者主动退出
再平衡协议:
- Join:消费者请求加入组
- Sync:Leader 消费者分配分区,同步给组员
- 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:
- 与 Leader 的延迟在配置时间内(
replica.lag.time.max.ms) - 落后消息数在阈值内(
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.factor | 3 | 副本因子,生产环境推荐 |
| min.insync.replicas | 2 | 最小同步副本数 |
| replica.lag.time.max.ms | 10000 | 同步超时时间 |
| unclean.leader.election.enable | false | 不允许非 ISR 选举 |
2.4.5 副本与容错
不同副本因子的容错能力
| 副本因子 | 可容忍故障数 | 说明 |
|---|---|---|
| 1 | 0 | 无容错,单点故障 |
| 2 | 0 | 允许 1 副本故障,但无法选举 |
| 3 | 1 | 允许 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 │
└──────────────────────────────┘ └──────────────────────────────┘
(每组独立消费) (每组独立消费)
下一步学习
- 第三章:配置详解 - 掌握 Broker、主题、生产者、消费者的配置
- 第四章:实现原理 - 深入理解 Kafka 内部机制
- 实践:运行你的第一个应用