Skip to content

Kafka

Kafka 是一个分布式事件流平台。它可以当消息队列用,但它的核心思想不是“消息被消费后立刻删除”,而是“把事件按顺序追加到日志里,消费者按自己的进度去读”。这也是 Kafka 和 RabbitMQ、RocketMQ 最大的思维差异。

零基础可以先这样理解:

RabbitMQ 更像快递分拣中心,消息投递给队列后由消费者取走;Kafka 更像一本不断追加的流水账,生产者负责往账本后面写,消费者记录自己读到哪一页。

Kafka 适合解决什么问题

Kafka 最常见的价值有三个:

  • 高吞吐:大量事件持续写入,例如日志、行为埋点、订单事件、设备数据。
  • 可回放:消息不会因为某个消费者读过就立刻消失,新消费者可以从旧位置重新读取。
  • 多订阅:同一份事件可以被风控、搜索、推荐、数据仓库等多个系统独立消费。
mermaid
flowchart TD
    A["业务系统产生事件"] --> B["Kafka Topic"]
    B --> C["实时风控"]
    B --> D["搜索索引"]
    B --> E["数据仓库"]
    B --> F["推荐系统"]

如果没有 Kafka,这些系统通常会互相同步调用,调用链会变长,峰值流量会互相拖垮,后续想补数也很困难。Kafka 把“事件产生”和“事件被谁使用”解耦,让业务系统只负责把事实写出去。

核心概念

概念通俗解释关键点
Record一条消息,也叫事件通常包含 key、value、timestamp、headers
Topic消息分类例如 order.createduser.behavior
PartitionTopic 的分片Kafka 扩展吞吐和保证局部顺序的核心
Offset分区内消息的位置消费者靠 offset 记录进度
BrokerKafka 服务节点负责存储分区、处理读写请求
Producer生产者把消息写入某个 Topic
Consumer消费者从 Topic 拉取消息处理
Consumer Group消费组同组内分摊消费,不同组独立消费
Replica副本用于容灾,一个分区可以有多个副本

Kafka 的整体工作流程

mermaid
flowchart TD
    A["Producer 构建事件"] --> B["根据 key 选择 Partition"]
    B --> C["写入 Partition Leader"]
    C --> D["Follower 复制数据"]
    D --> E["达到确认条件后返回成功"]
    E --> F["Consumer Group 拉取消息"]
    F --> G["业务处理成功"]
    G --> H["提交 Offset"]

这个流程里最重要的是三个问题:

  • 写到哪个分区:决定顺序性和并发度。
  • 写成功的标准是什么:决定消息可靠性。
  • 消费到哪里了:决定重复消费和漏消费风险。

为什么 Kafka 用追加日志

Kafka 不像传统队列那样频繁在队列头部删除消息,而是把消息顺序追加到文件末尾。这样设计有几个好处:

  • 顺序写磁盘比随机写快,吞吐更稳定。
  • 消费者只保存 offset,不需要 Broker 为每个消费者复制一份消息。
  • 消息可以保留一段时间,支持回放、补数、重新建索引。
  • 多个消费组互不影响,同一份数据可以服务多个业务。
mermaid
flowchart TD
    A["Partition 日志文件"] --> B["Offset 0"]
    B --> C["Offset 1"]
    C --> D["Offset 2"]
    D --> E["Offset 3"]
    F["消费组 A 读到 Offset 2"] --> D
    G["消费组 B 读到 Offset 0"] --> B

如果你把 Kafka 只当“普通队列”使用,就容易忽略 offset、分区、保留时间、重复消费这些问题。结果通常是:消息能跑,但一到重启、扩容、补数据、故障恢复,就开始丢、重、乱。

学习路线

Kafka、RabbitMQ、RocketMQ 怎么选

对比项RabbitMQKafkaRocketMQ
核心模型Exchange + QueueTopic + Partition + LogTopic + Queue + ConsumerGroup
典型定位业务消息、复杂路由、任务队列事件流、日志、数据管道、高吞吐业务消息、事务消息、顺序消息
消息保留偏队列消费,消费后可删除按时间或大小保留,可回放支持堆积和重试,偏业务消息
顺序能力单队列可顺序分区内顺序队列内顺序
路由能力很强,Exchange 类型丰富弱,主要靠 Topic 和 key中等,靠 Topic、Tag、Key
适合场景企业系统解耦、复杂路由、延迟/死信日志、埋点、数据同步、实时计算交易链路、事务消息、消费重试

如果你只是做传统业务异步通知,RabbitMQ 上手更直观;如果你要承接大量事件并允许多个系统重复读取,Kafka 更合适;如果你特别关注事务消息、订单链路和业务重试,RocketMQ 往往更贴近业务语义。

最小业务 Demo:订单事件

生产者发送的是“事实”,不是“命令”。推荐用过去式命名,例如 OrderCreatedEvent,表示订单已经创建。

java
public record OrderCreatedEvent(Long orderId, Long userId, Long amount) {
}
java
@Service
public class OrderEventProducer {
    private final KafkaTemplate<String, OrderCreatedEvent> kafkaTemplate;

    public OrderEventProducer(KafkaTemplate<String, OrderCreatedEvent> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void send(OrderCreatedEvent event) {
        kafkaTemplate.send("order.created", String.valueOf(event.orderId()), event);
    }
}

消费端一定要做幂等,因为 Kafka 常见可靠语义是“至少一次”,业务处理可能重复执行。

java
@Component
public class OrderCreatedConsumer {

    @KafkaListener(topics = "order.created", groupId = "search-index-service")
    public void onMessage(OrderCreatedEvent event) {
        if (hasIndexed(event.orderId())) {
            return;
        }
        rebuildOrderIndex(event.orderId());
        markIndexed(event.orderId());
    }

    private boolean hasIndexed(Long orderId) {
        return false;
    }

    private void rebuildOrderIndex(Long orderId) {
        System.out.println("同步订单到搜索索引:" + orderId);
    }

    private void markIndexed(Long orderId) {
    }
}

这个 Demo 背后的原则是:Kafka 负责保存事件和推进消费进度,业务系统负责保证事件处理幂等。两者缺一不可。

官方资料