Skip to content

Kafka

这篇笔记聚焦 Kafka 这个具体产品,重点放在 TopicPartitionOffsetConsumer Group吞吐顺序回放Spring Boot 集成 这些最常和 Kafka 一起出现的主线问题上。

如果你还没建立消息队列这条知识线,可以看:

这篇笔记不会把 Kafka 写成“另一种普通队列”,而是重点强调:Kafka 更像一个高吞吐、可保留、可回放的事件日志流平台。


1. Kafka 是什么

Kafka 可以看成 一个面向分布式日志流和事件流的消息平台

这里的重点不是“发一条消息”,而是:

  1. 持续承接大量事件写入
  2. 把同一份数据分发给多个下游系统
  3. 允许不同消费者按自己的进度消费
  4. 在一段时间内保留历史消息,支持回放

所以 Kafka 更常被用在这些语境里:

  1. 日志采集
  2. 用户行为埋点
  3. 订单、支付、风控等业务事件流
  4. 实时计算和流式处理前置链路

它不是简单拿来和 RabbitMQ 做“一换一”的替代,而是:更适合对吞吐、日志保留、分区扩展和事件回放更敏感的场景。


2. 为什么 Kafka 常和日志流、事件流放在一起讨论

很多业务消息系统更像:来了一个任务,尽快投给下游处理。

Kafka 的思路更像:把一连串持续产生的事件,按顺序不断写进日志,再让不同消费者按各自节奏读取。

这会带来几个很重要的差异:

  1. 消息不会因为某个消费者读过就立刻消失
  2. 不同消费者组可以独立消费同一份数据
  3. 消费进度是“消费者自己记到哪里了”
  4. 历史数据在保留期内还可以重放

也正因为这样,Kafka 天然适合:

  1. 一份数据被多个系统复用
  2. 需要回溯、补算、重放历史消息
  3. 需要持续处理大规模数据流

3. Kafka 里最核心的几个概念

3.1 Broker

Broker 可以理解成 真正负责接收、存储和分发消息的 Kafka 服务节点

如果是本地学习环境,常见地址就是:

  1. localhost:9092

3.2 Topic

Topic 可以理解成 同一类消息的逻辑主题

例如:

  1. order-events
  2. payment-events
  3. user-behavior
  4. audit-logs

3.3 Partition

Partition 是 Kafka 最核心的概念之一。

它可以理解成 一个 Topic 在物理上拆成的多个分区

它主要解决两件事:

  1. 提升并行处理能力
  2. 让同一个业务 key 可以稳定落到同一个分区里,从而维持局部顺序

3.4 Offset

Offset 可以理解成 消息在某个分区中的位置编号

它表达的是:

  1. 这条消息在当前分区里排第几个
  2. 消费者已经处理到哪里

要注意的是:Kafka 更常是“消息保留在日志里,消费者自己记录进度”,而不是“消息被消费后就立刻删掉”。

3.5 Producer

Producer 就是消息生产者。

谁真正把消息写进 Kafka,谁就是 Producer

它可能是:

  1. 订单服务
  2. 支付服务
  3. 日志接入服务
  4. 采集日志的 Agent

3.6 Consumer

Consumer 就是消息消费者。

谁从 Kafka 拉取消息并处理,谁就是 Consumer

它的处理动作常包括:

  1. 落库
  2. 发告警
  3. 做清洗转换
  4. 写搜索系统
  5. 转发给别的下游

3.7 Consumer Group

Consumer Group 可以理解成 一组协同消费同一个 Topic 的消费者实例

它最值得记住的规则是:

  1. 同一个组内,同一个分区同一时刻只会分给一个消费者
  2. 不同组之间互不影响,可以各自独立消费同一份数据

这也是 Kafka 为什么特别适合 一份事件,多个系统,各自独立处理


4. 一条 Kafka 消息通常是怎么流转的

看这张整体图:

mermaid
flowchart LR
    P[Producer 发送消息] --> T[Topic]
    T --> P0[Partition 0]
    T --> P1[Partition 1]
    T --> P2[Partition 2]
    P0 --> G1[Consumer Group A]
    P1 --> G1
    P2 --> G1
    P0 --> G2[Consumer Group B]
    P1 --> G2
    P2 --> G2

这条链路真正表达的是:

  1. 生产者不是把消息直接发给某个消费者
  2. 而是先写进 Topic
  3. Topic 再按分区存储消息
  4. 不同消费者组各自维护自己的消费进度

所以同一份 order-events 数据,可以同时被:

  1. 订单分析服务消费
  2. 审计服务消费
  3. 风控服务消费
  4. 搜索索引服务消费

5. Kafka 的顺序、吞吐和回放应该怎么理解

5.1 顺序:更常保证分区内顺序

Kafka 经常被说“支持顺序消息”,但更贴近实际的是:Kafka 更常保证分区内顺序,而不是全局顺序。

也就是说:

  1. 同一个分区里的消息天然有序
  2. 不同分区之间不保证全局有序

所以如果业务要求:同一个订单的事件必须按顺序处理

更常见的做法是:让同一个 orderId 作为 message key,这样它更稳定地落到同一个分区。

5.2 吞吐:为什么 Kafka 通常更高

Kafka 的高吞吐,通常和这些设计有关:

  1. 顺序追加写入
  2. 分区带来的并行扩展能力
  3. 批量发送和批量拉取
  4. 更偏日志流的模型设计

它不是靠某一个技巧“突然变快”,而是:存储模型、分区模型和批处理模型一起服务于高吞吐。

5.3 回放:为什么 Kafka 特别适合补消费

Kafka 很重要的一点是:消息在保留期内通常还在,消费者只是记录自己读到哪里。

所以当出现下面这些场景时,它会比较自然:

  1. 新增一个消费者组,重新消费历史消息
  2. 某个消费者故障恢复后,从上次位点继续
  3. 因为业务 bug 需要补算历史数据

这也是它和很多“处理完就结束”的业务队列非常不一样的地方。


6. Kafka 更适合哪些场景,不适合哪些场景

更适合

  1. 日志采集
  2. 埋点数据流
  3. 大规模业务事件流
  4. 实时计算前置链路
  5. 一份数据被多个下游独立消费

不那么适合

  1. 特别依赖复杂路由规则的业务消息
  2. 更强调交换器模型和细粒度投递控制的场景
  3. 需要把“一条消息怎么分发到多个队列”表达得非常细的场景

如果业务重点是:消息怎么按复杂路由规则投给多个下游

RabbitMQ 往往更自然。


7. 企业里常见的 Kafka 使用形态

结合实际项目,Kafka 常见的几种形态可以压缩成下面三类。

7.1 业务系统直接写 Kafka

这种方式最常见于事件驱动架构:

  1. 订单服务写 order-events
  2. 支付服务写 payment-events
  3. 下游审计、风控、报表系统各自订阅

它的优点是:

  1. 解耦清晰
  2. 一份事件可以被多个系统复用
  3. 扩展性较好

7.2 日志文件或标准输出先被 Agent 采集,再发 Kafka

这种方式更常见于日志平台:

  1. 应用先写本地日志
  2. FilebeatFluent BitVector 等 Agent 采集
  3. 再统一写入 Kafka

它更适合:

  1. 老系统较多
  2. 多语言混合环境
  3. 日志以文件形态存在

7.3 Kafka 前面接入事件,后面再接处理层和落库层

这种方式更像成熟企业架构:

  1. 上游系统把原始事件写进 Kafka
  2. 中间用 FlinkKafka Streams 或自定义服务做清洗和加工
  3. 最后再落到 MySQLElasticsearchClickHouse 等系统

它更适合:

  1. 实时报表
  2. 实时风控
  3. 实时告警
  4. 大规模日志分析

8. Spring Boot 项目里怎么整合 Kafka

如果是 Spring Boot 项目,最常见的目标通常是:

  1. 项目能连接 Kafka
  2. 能声明或初始化 Topic
  3. 能写生产者
  4. 能写消费者
  5. 能明确消息模型和消费组配置

8.1 依赖通常长什么样

最基础的一组依赖通常包括:

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

    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-validation</artifactId>
    </dependency>

    <dependency>
        <groupId>com.fasterxml.jackson.datatype</groupId>
        <artifactId>jackson-datatype-jsr310</artifactId>
    </dependency>

    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <optional>true</optional>
    </dependency>
</dependencies>

这几项分别在解决:

  1. spring-kafka:Spring Boot 和 Kafka 的核心集成能力
  2. validation:校验自定义配置类
  3. jackson-datatype-jsr310:更自然地处理 LocalDateTimeInstant 等时间类型
  4. lombok:减少样板代码

8.2 application.yml 通常怎么配

一个更贴近 Spring Boot 工程的基础配置通常会是这样:

yml
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
    consumer:
      group-id: order-analytics-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.demo.messaging
            value:
              default:
                type: com.example.demo.messaging.OrderPaidEvent
    listener:
      ack-mode: record

app:
  kafka:
    topic:
      order-events: order.events
      partitions: 3
      replicas: 1

这里最值得记住的几个配置是:

  1. bootstrap-servers:Kafka Broker 地址
  2. producer:生产者序列化和确认策略
  3. consumer.group-id:当前消费者属于哪个消费组
  4. auto-offset-reset:没有已提交位点时,从哪里开始消费
  5. listener.ack-mode:监听容器按什么粒度确认消费

8.3 Topic 声明代码通常长什么样

在 Spring Boot 项目里,常见做法是用 NewTopic 在启动阶段声明 Topic 元数据:

java
@Configuration
public class KafkaTopicConfiguration {

    /**
     * orderEventsTopic
     * 功能:声明订单事件 Topic。
     * 核心职责:
     * 1. 在应用启动时向 KafkaAdmin 注册 Topic 元数据
     * 2. 统一管理 Topic 名称、分区数和副本数
     * 参数:
     * - kafkaTopicProperties: 自定义 Kafka Topic 配置
     * 返回值:
     * - NewTopic: 订单事件 Topic 定义
     */
    @Bean
    public NewTopic orderEventsTopic(KafkaTopicProperties kafkaTopicProperties) {
        return TopicBuilder
            .name(kafkaTopicProperties.getOrderEvents())
            .partitions(kafkaTopicProperties.getPartitions())
            .replicas(kafkaTopicProperties.getReplicas())
            .build();
    }
}

这段代码真正表达的是:Topic 不只是一个字符串名字,它还包含分区数、副本数这些很关键的元数据。

8.4 消息模型通常长什么样

如果项目走 JSON 消息体,通常会有一个明确的事件模型:

java
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class OrderPaidEvent {

    /**
     * 事件唯一标识,用于区分每一条订单支付事件。
     */
    private String eventId;

    /**
     * 订单号。工程上常常也会把它作为 Kafka message key。
     */
    private String orderId;

    /**
     * 用户 ID,用于标识当前事件属于哪个用户。
     */
    private Long userId;

    /**
     * 支付金额。
     */
    private BigDecimal amount;

    /**
     * 事件发生时间。
     */
    private LocalDateTime paidAt;
}

这里把 orderId 单独拎出来很重要,因为它经常承担两层角色:

  1. 消息体里的业务字段
  2. message key

第二层角色会直接影响分区路由和局部顺序。

8.5 生产者侧通常长什么样

Spring Boot 项目里最常见的生产者写法,是通过 KafkaTemplate 发送事件:

java
@Service
@RequiredArgsConstructor
public class OrderEventProducer {

    private final KafkaTemplate<String, OrderPaidEvent> kafkaTemplate;
    private final KafkaTopicProperties kafkaTopicProperties;

    /**
     * sendOrderPaidEvent
     * 功能:发送订单支付成功事件。
     * 核心职责:
     * 1. 组装业务事件对象
     * 2. 指定 Topic 和 message key
     * 3. 把支付成功这件事发布到 Kafka
     * 参数:
     * - orderId: 订单号
     * - userId: 用户 ID
     * - amount: 支付金额
     * 返回值:
     * - void
     */
    public void sendOrderPaidEvent(String orderId, Long userId, BigDecimal amount) {
        OrderPaidEvent event = OrderPaidEvent.builder()
            .eventId(UUID.randomUUID().toString())
            .orderId(orderId)
            .userId(userId)
            .amount(amount)
            .paidAt(LocalDateTime.now())
            .build();

        // 让同一个订单号稳定作为 message key,有助于维持局部顺序。
        kafkaTemplate.send(kafkaTopicProperties.getOrderEvents(), orderId, event);
    }
}

这段生产者代码真正想说明的,不只是“怎么发消息”,而是:

  1. 生产者通常放在应用服务层
  2. 业务事件要有明确模型
  3. message key 往往要和业务顺序要求绑在一起考虑

8.6 消费者侧通常长什么样

消费者最常见的写法,是用 @KafkaListener 监听 Topic:

java
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderAnalyticsConsumer {

    private final OrderAnalyticsService orderAnalyticsService;

    /**
     * handleOrderPaid
     * 功能:消费订单支付成功事件,并写入分析系统或下游存储。
     * 核心职责:
     * 1. 监听订单事件 Topic
     * 2. 获取消息体以及分区、位点等元数据
     * 3. 完成业务处理后返回,由监听容器按配置确认消费
     * 参数:
     * - event: 反序列化后的订单支付事件
     * - partition: 当前消息所在分区
     * - offset: 当前消息的位点
     * - messageKey: 当前消息的 key
     * 返回值:
     * - void
     */
    @KafkaListener(
        topics = "${app.kafka.topic.order-events}",
        groupId = "order-analytics-group",
        concurrency = "3"
    )
    public void handleOrderPaid(
            OrderPaidEvent event,
            @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
            @Header(KafkaHeaders.OFFSET) long offset,
            @Header(value = KafkaHeaders.RECEIVED_KEY, required = false) String messageKey) {

        log.info(
            "收到订单支付事件 orderId={}, key={}, partition={}, offset={}",
            event.getOrderId(),
            messageKey,
            partition,
            offset
        );

        orderAnalyticsService.recordPaidEvent(event);
    }
}

这里最值得注意的几个点是:

  1. topics:监听哪个 Topic
  2. groupId:当前消费者属于哪个消费组
  3. concurrency = "3":监听容器允许的并发消费者数量
  4. partition / offset / key:这些元数据对排查顺序、定位消费进度很有帮助

8.7 这些 Spring Boot 代码到底表达了什么

把前面的 Topic、Producer、Consumer 串起来看,完整链路其实是:

mermaid
flowchart LR
    A[订单服务完成支付] --> B[OrderEventProducer]
    B --> C[Kafka order.events Topic]
    C --> D[OrderAnalyticsConsumer]
    D --> E[OrderAnalyticsService]

这条链路真正表达的是:

  1. 业务系统负责把“发生了什么事”发布成事件
  2. Kafka 负责承接和分发这份事件
  3. 消费者组按自己的节奏处理消息
  4. 下游服务再决定落库、统计、告警还是继续转发

9. Kafka 集成里最常见的工程关注点

9.1 分区数不是越多越好

分区数会影响:

  1. 并行消费上限
  2. 顺序粒度
  3. 运维和管理复杂度

所以它不是越大越高级,而是要看:业务吞吐、消费者并发和顺序要求到底是什么。

9.2 message key 设计会直接影响顺序

如果业务要求同一个订单、有同一个用户或同一个设备的消息尽量有序,通常要:

  1. 明确业务 key
  2. 让这个 key 参与分区路由

如果不认真设计 message key,很容易出现:

  1. 消息被均匀打散到多个分区
  2. 吞吐上去了
  3. 但同一业务主体的顺序丢了

9.3 消费成功不等于业务完全没风险

即使 @KafkaListener 收到了消息,也不等于所有问题都解决了。

还要继续考虑:

  1. 业务处理是否幂等
  2. 异常时怎么重试
  3. 消费位点什么时候提交更合适
  4. 下游存储失败时怎么补偿

Kafka 很强,但它不会自动替业务解决所有一致性问题。


10. 一句话总结

Kafka 最值得记住的是:

它更像一个高吞吐、可分区、可保留、可回放的事件日志流平台。放到 Spring Boot 项目里,最常见的落地方式就是:用 Topic 承接事件、用 KafkaTemplate 发送消息、用 @KafkaListener 消费消息,再围绕分区、位点、消费组和 message key 去理解吞吐、顺序和工程边界。

基于 VitePress 构建的个人技术笔记。