Kafka 我之前在项目里断断续续接触过,但一直没把它和 RabbitMQ 的定位差异、以及"它为什么是日志流而不是普通消息队列"这件事想透。这篇把我梳理的结果记下来:核心概念、和 RabbitMQ 的区别、消息怎么不丢、顺序怎么保证、怎么和 Spring Boot 接上。
一、Kafka 是什么,和 RabbitMQ 有什么不一样
两个都是消息中间件,但设计哲学完全不同,用错场景会很别扭:
| 维度 | Kafka | RabbitMQ |
|---|---|---|
| 本质 | 分布式日志流(append-only 日志) | 消息路由器(AMQP 协议,路由规则) |
| 核心模型 | Topic → Partition(分区) | Exchange → Binding → Queue |
| 吞吐 | 极高(顺序写磁盘,百万级/秒) | 中(万级/秒) |
| 消息消费 | 拉取(consumer 主动 pull) | 推送(broker push 给 consumer) |
| 消费后 | 消息保留(可重复消费、回放) | 消费确认后删除 |
| 典型场景 | 日志采集、流处理、事件溯源、大数据管道 | 业务解耦、任务队列、RPC、延迟队列 |
一句话记住:RabbitMQ 是"投递完就忘"的消息路由,Kafka 是"可以反复读的历史日志"。
这决定了选型——你要"一个消息准确投给一个消费者、处理完就删",用 RabbitMQ;你要"海量事件流、多个消费者各读各的、还能回溯历史",用 Kafka。更细的对比可以看《RabbitMQ 快速入门》。
二、核心概念
Kafka 的概念是一条链,顺着数据流向理解:
Producer → Topic → Partition → Consumer(Consumer Group)
│
└─ 每个 Partition 是一个有序、不可变的日志
offset 是消息在分区内的唯一位置| 概念 | 作用 | 类比 |
|---|---|---|
| Broker | Kafka 服务节点,一个集群多个 Broker | 一台装日志的服务器 |
| Topic | 消息的逻辑分类(如 order-events) | 一个"主题"文件夹 |
| Partition | Topic 的物理分片,并行 + 有序的最小单位 | 文件夹里的多个日志文件 |
| Offset | 消息在分区内的位置编号(递增) | 日志的行号 |
| Producer | 发消息,指定 Topic + 分区 | 写日志的人 |
| Consumer Group | 一组消费者,组内分摊分区、组间互不影响 | 一群分工读日志的人 |
两个最容易混的点:
分区 = 并行度 + 顺序边界:一个 Topic 分几个分区,就能被几个消费者并行消费;但顺序只在单个分区内保证,跨分区无序。所以"要严格顺序"的消息,必须保证发到同一个分区(比如用同一个 key)。
消费组:同组内的消费者瓜分分区(一个分区只能被组内一个消费者消费,避免重复);不同组的消费者各自独立消费同一份消息(互不影响)。这就是 Kafka"一份消息多路消费"的来源——比如同一个订单事件,订单组和风控组各消费一份。
三、消息怎么不丢
Kafka 的可靠性靠三个环节配合,任何一个没配好都可能丢:
① 生产者 acks:
acks=all # 所有 ISR 副本都确认才返回成功(最可靠)
# acks=1 # leader 确认即可(可能丢,leader 挂了还没同步)
# acks=0 # 不确认,最快但最容易丢② 副本(Replication):每个分区有多个副本(leader + follower),leader 挂了 follower 顶上。replication.factor 决定副本数,min.insync.replicas 决定最少几个副本同步成功才算写成功。
③ 消费者 offset 提交:消费者处理完消息要提交 offset,否则重启后会重复消费。关键在提交时机——先处理再提交(enable.auto.commit=false + 手动提交),避免"先提交了但业务没处理完,进程挂了就丢这条"。
四、Spring Boot 集成
1. 依赖:
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>2. 生产者(KafkaTemplate):
@Service
public class OrderProducer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void send(String orderId) {
kafkaTemplate.send("order-events", orderId, "{\"orderId\":\"" + orderId + "\"}");
}
}3. 消费者(@KafkaListener):
@Component
public class OrderConsumer {
@KafkaListener(topics = "order-events", groupId = "order-group")
public void onMessage(String message) {
// 处理订单事件
System.out.println("收到: " + message);
}
}4. 关键配置(application.yml):
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
acks: all # 可靠性:所有副本确认
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
consumer:
group-id: order-group
enable-auto-commit: false # 手动提交 offset,避免丢消息
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer五、Docker 快速起一个(本地验证)
Kafka 依赖 ZooKeeper(新版可不用,这里用最经典的组合):
# 1. 起 ZooKeeper
docker run -d --name zookeeper -p 2181:2181 zookeeper:3.8
# 2. 起 Kafka
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
--link zookeeper \
apache/kafka:3.7.0起好后用自带的脚本验证收发:
# 创建 topic
docker exec kafka kafka-topics.sh --create --topic test --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1
# 发一条消息
echo "hello kafka" | docker exec -i kafka kafka-console-producer.sh --topic test --bootstrap-server localhost:9092
# 收消息
docker exec kafka kafka-console-consumer.sh --topic test --from-beginning --bootstrap-server localhost:9092小结
- Kafka 是分布式日志流(append-only,可回放),不是 RabbitMQ 那种"消费即删"的消息路由
- 核心:Topic 逻辑分类 → Partition 物理分片(并行 + 顺序的最小单位)→ Consumer Group 组内分摊组间独立
- 顺序只在分区内保证,跨分区无序
- 不丢消息三件事:
acks=all+ 多副本 + 手动提交 offset - Spring Boot 用
KafkaTemplate发、@KafkaListener收,配好序列化和 acks
下一步想了解 Kafka 和 RabbitMQ 到底怎么选,可以看《RabbitMQ 快速入门》里的对比,或者深入研究 Kafka 的分区策略、重复消费幂等、消息积压这些生产级问题(文章整理中)。
