第 13 章:Kafka 进阶
学习目标
- 理解 Kafka 的分区、副本、ISR 机制
- 掌握生产者/消费者 Java 代码
- 处理消息顺序、幂等、事务
- 监控 Kafka 集群状态
一、Kafka 核心概念
| 概念 | 含义 |
|---|---|
| Producer | 消息生产者 |
| Consumer | 消息消费者 |
| Broker | 一个 Kafka 进程 |
| Topic | 消息主题 |
| Partition | 分区,物理存储单位 |
| Offset | 消息在分区内的编号 |
| Consumer Group | 消费者组,组内均分分区 |
| ISR | In-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
# 看消费进度、lagLag 持续增大 → 消费慢,要扩容消费者或优化业务。
十一、本章小结
| 要点 | 关键 |
|---|---|
| 分区 | topic 物理切分,同分区有序 |
| acks=all + 幂等 | 生产侧不丢不重 |
| 消费者组 | 同组内均分分区 |
| offset | earliest 防止新组丢消息 |
| 幂等 | Redis SETNX 判重 |
| 事务 | 跨 topic / 跨分区原子 |
| lag | 监控消费延迟,及时扩容 |
动手练习
- 用 Docker 启动 Kafka(单 broker + 3 分区 topic)
- 生产者用 userId 当 key,确保同一用户的消息落到同分区
- 启动两个消费者(同一 group),观察分区被均分
下一章:第 14 章:RocketMQ →