Skip to content

Kafka 我之前在项目里断断续续接触过,但一直没把它和 RabbitMQ 的定位差异、以及"它为什么是日志流而不是普通消息队列"这件事想透。这篇把我梳理的结果记下来:核心概念、和 RabbitMQ 的区别、消息怎么不丢、顺序怎么保证、怎么和 Spring Boot 接上。

一、Kafka 是什么,和 RabbitMQ 有什么不一样

两个都是消息中间件,但设计哲学完全不同,用错场景会很别扭:

维度KafkaRabbitMQ
本质分布式日志流(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 是消息在分区内的唯一位置
概念作用类比
BrokerKafka 服务节点,一个集群多个 Broker一台装日志的服务器
Topic消息的逻辑分类(如 order-events一个"主题"文件夹
PartitionTopic 的物理分片,并行 + 有序的最小单位文件夹里的多个日志文件
Offset消息在分区内的位置编号(递增)日志的行号
Producer发消息,指定 Topic + 分区写日志的人
Consumer Group一组消费者,组内分摊分区、组间互不影响一群分工读日志的人

两个最容易混的点

  1. 分区 = 并行度 + 顺序边界:一个 Topic 分几个分区,就能被几个消费者并行消费;但顺序只在单个分区内保证,跨分区无序。所以"要严格顺序"的消息,必须保证发到同一个分区(比如用同一个 key)。

  2. 消费组:同组内的消费者瓜分分区(一个分区只能被组内一个消费者消费,避免重复);不同组的消费者各自独立消费同一份消息(互不影响)。这就是 Kafka"一份消息多路消费"的来源——比如同一个订单事件,订单组和风控组各消费一份。

三、消息怎么不丢

Kafka 的可靠性靠三个环节配合,任何一个没配好都可能丢:

① 生产者 acks

properties
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. 依赖

xml
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

2. 生产者KafkaTemplate):

java
@Service
public class OrderProducer {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    public void send(String orderId) {
        kafkaTemplate.send("order-events", orderId, "{\"orderId\":\"" + orderId + "\"}");
    }
}

3. 消费者@KafkaListener):

java
@Component
public class OrderConsumer {
    @KafkaListener(topics = "order-events", groupId = "order-group")
    public void onMessage(String message) {
        // 处理订单事件
        System.out.println("收到: " + message);
    }
}

4. 关键配置(application.yml):

yaml
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(新版可不用,这里用最经典的组合):

bash
# 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

起好后用自带的脚本验证收发:

bash
# 创建 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 的分区策略、重复消费幂等、消息积压这些生产级问题(文章整理中)。