第 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 32.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: 1bash
docker-compose up -d4.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 | 不等确认 | 可能丢 |
| 1 | Leader 写入即返回 | 极少数丢 |
| 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_immediate6.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 + 事务 | 否 | 慢 |
动手练习
- 用 Docker Compose 启动 Kafka
- 用 Spring Boot 实现 Producer / Consumer
- 手动提交 offset 实现至少一次
- 用 key 保证同一订单顺序
推荐阅读
下一章:第 199 章:Kafka 进阶与可靠性 →