Kafka 异步消息推送与事件驱动架构详解


一、为什么需要 Kafka

先从一个实际问题说起:

假设有一个电商系统,用户下单后需要:

  1. 扣减库存
  2. 通知仓库发货
  3. 发送短信
  4. 更新积分
  5. 同步数据到第三方系统

同步调用的问题:

用户下单 → 扣库存(50ms) → 通知仓库(200ms) → 发短信(300ms) → 更新积分(100ms) → 同步第三方(500ms)
总耗时:1150ms,用户等待体验极差

问题:

  • 耗时叠加,用户等太久
  • 任何一个环节失败,整个下单失败
  • 新增一个通知渠道,就要改下单代码
  • 下游系统挂了,上游也跟着挂

Kafka 解决方案:

用户下单 → 扣库存(50ms) → 发送"订单已创建"事件到Kafka(5ms) → 返回成功
                                        ↓
                            Kafka 异步分发给各消费者:
                            ├── 仓库服务消费 → 通知发货
                            ├── 短信服务消费 → 发送短信
                            ├── 积分服务消费 → 更新积分
                            └── 第三方同步消费 → 同步数据
总耗时:55ms,用户几乎无感

注:

博客:

https://blog.csdn.net/badao_liumang_qizhi

二、Kafka 核心概念

2.1 整体架构

┌────────────────────────────────────────────────────────────────┐
│                        Kafka Cluster                           │
│                                                                │
│  ┌──────────┐    ┌──────────┐    ┌──────────┐                │
│  │ Broker 0 │    │ Broker 1 │    │ Broker 2 │                │
│  │          │    │          │    │          │                │
│  │ Topic-A  │    │ Topic-A  │    │ Topic-A  │                │
│  │ Part-0   │    │ Part-1   │    │ Part-2   │                │
│  └──────────┘    └──────────┘    └──────────┘                │
│                                                                │
└────────────────────────────────────────────────────────────────┘
        ↑                                          ↓
   ┌─────────┐                              ┌──────────┐
   │Producer │  发送消息                     │Consumer  │ 拉取消息
   │(生产者) │ ─────────→ Kafka ──────────→ │(消费者)  │
   └─────────┘                              └──────────┘

2.2 核心术语

Broker(代理/节点)

Kafka 集群中的一台服务器就是一个 Broker。多个 Broker 组成 Kafka 集群,提供高可用和水平扩展能力。每个 Broker 用唯一的 ID 标识。

类比:Broker 就像邮局的一个网点,多个网点组成整个邮政系统

Topic(主题)

消息的逻辑分类。Producer 将消息发送到指定 Topic,Consumer 从指定 Topic 消费消息。一个 Topic 可以有多个 Producer 和多个 Consumer。

类比:Topic 就像报纸的版面(体育版、财经版),读者订阅自己感兴趣的版面

Partition(分区)

一个 Topic 可以分为多个 Partition。Partition 是 Kafka 并行处理的基本单位。每个 Partition 内的消息是有序的,但跨 Partition 不保证顺序。

类比:一条高速公路(Topic)有多个车道(Partition),每个车道内车是有序的,
      但不同车道的车相对顺序不确定
Topic: order-events (3个分区)
├── Partition 0: [msg-0, msg-3, msg-6, msg-9 ...]
├── Partition 1: [msg-1, msg-4, msg-7, msg-10 ...]
└── Partition 2: [msg-2, msg-5, msg-8, msg-11 ...]

Offset(偏移量)

每条消息在 Partition 内的唯一序号,从 0 开始递增。Consumer 通过 Offset 记录自己消费到了哪里,实现断点续读。

Partition 0: [msg@offset0, msg@offset1, msg@offset2, msg@offset3 ...]
                                              ↑
                                    Consumer当前消费位置

Producer(生产者)

负责将消息发送到 Kafka 的 Topic 中。Producer 可以选择将消息发送到哪个 Partition(通过 key 哈希、轮询、或自定义策略)。

Consumer(消费者)

负责从 Kafka 的 Topic 中拉取消息并处理。Consumer 主动拉取(Pull 模式),而不是 Kafka 推送。

Consumer Group(消费者组)

多个 Consumer 组成的组。同一个 Group 内的 Consumer 分摊消费 Topic 中的 Partition(负载均衡)。不同 Group 各自独立消费全量消息(广播)。

Topic: order-events (3个分区)

Consumer Group A (订单服务):
├── Consumer-A1 消费 Partition 0
├── Consumer-A2 消费 Partition 1
└── Consumer-A3 消费 Partition 2

Consumer Group B (通知服务):
├── Consumer-B1 消费 Partition 0, 1
└── Consumer-B2 消费 Partition 2

两个Group各自独立消费全量数据

Replica(副本)

每个 Partition 可以有多个副本,分布在不同 Broker 上。一个 Leader 负责读写,多个 Follower 负责同步数据。Leader 挂了,Follower 自动选举为新 Leader。

2.3 消息流转全流程

1. Producer 发送消息
   ┌─────────────────────────────────────────────────┐
   │ Producer                                         │
   │   → 序列化消息(对象→字节数组)                  │
   │   → 选择 Partition(根据key hash或轮询)         │
   │   → 消息放入发送缓冲区(RecordAccumulator)      │
   │   → Sender 线程批量发送到 Broker                 │
   │   → 等待 ack 确认(可配置)                      │
   └─────────────────────────────────────────────────┘

2. Broker 存储消息
   ┌─────────────────────────────────────────────────┐
   │ Broker                                           │
   │   → 接收消息,写入对应 Partition 的日志文件       │
   │   → 分配 Offset                                  │
   │   → 同步到 Follower 副本                         │
   │   → 返回 ack 给 Producer                        │
   └─────────────────────────────────────────────────┘

3. Consumer 消费消息
   ┌─────────────────────────────────────────────────┐
   │ Consumer                                         │
   │   → 向 Broker 发送 fetch 请求(带 offset)       │
   │   → 拉取一批消息                                 │
   │   → 反序列化(字节数组→对象)                    │
   │   → 业务处理                                     │
   │   → 提交 offset(记录消费位置)                  │
   └─────────────────────────────────────────────────┘

三、Kafka 关键 API

3.1 Producer API

API / 配置 说明
send(topic, key, value) 发送消息,key 决定分区路由
send(topic, value) 发送消息,轮询分区
flush() 强制将缓冲区消息发送出去
close() 关闭 Producer,释放资源
acks=0 不等待确认,最快但可能丢消息
acks=1 Leader 写入即确认,平衡性能和可靠性
acks=all 所有副本写入才确认,最安全但最慢
retries 发送失败重试次数
batch.size 批量发送的大小阈值
linger.ms 等待更多消息凑批的最大时间

3.2 Consumer API

API / 配置 说明
subscribe(topics) 订阅 Topic 列表
poll(timeout) 拉取消息(核心循环)
commitSync() 同步提交 offset
commitAsync() 异步提交 offset
seek(partition, offset) 指定从某个 offset 开始消费
seekToBeginning() 从头消费
seekToEnd() 从最新位置消费
group.id 消费者组标识
auto.offset.reset 无 offset 时策略:earliest/latest
enable.auto.commit 是否自动提交 offset
max.poll.records 每次 poll 最多拉取条数

3.3 消息可靠性保证级别

级别 含义 适用场景
At Most Once 最多一次,可能丢消息 日志采集、监控数据
At Least Once 至少一次,可能重复 大多数业务场景(配合幂等)
Exactly Once 精确一次,不丢不重 金融转账、对账(代价最高)

四、Kafka 与事件驱动架构的关系

4.1 事件驱动架构回顾

传统同步调用:
A服务 → 直接调用 → B服务 → 直接调用 → C服务
        (强耦合,A要知道B的存在)

事件驱动:
A服务 → 发布事件 → [Kafka] ← 订阅事件 ← B服务
                            ← 订阅事件 ← C服务
                            ← 订阅事件 ← D服务(新增消费者,A无感知)
        (松耦合,A不需要知道谁在消费)

4.2 Kafka 在事件驱动中的角色

┌───────────┐     事件      ┌─────────────────┐    事件     ┌───────────┐
│ 事件生产者 │ ──────────→  │   Kafka         │ ─────────→ │ 事件消费者 │
│           │              │ (事件通道/总线)   │            │           │
│ - 订单服务 │              │                  │            │ - 库存服务 │
│ - 支付服务 │              │ 职责:           │            │ - 物流服务 │
│ - 用户服务 │              │ · 持久化存储事件  │            │ - 通知服务 │
└───────────┘              │ · 保证投递可靠性  │            │ - 分析服务 │
                           │ · 支持回溯重放    │            └───────────┘
                           │ · 解耦生产消费    │
                           └─────────────────┘

4.3 核心价值

价值 说明
时间解耦 生产者和消费者不需要同时在线
空间解耦 生产者不需要知道消费者的地址
同步解耦 生产者发送后立即返回,不等待消费者处理完成
流量缓冲 突发流量被 Kafka 缓存,消费者按自己速度处理
事件溯源 消息持久化,可以回溯历史事件,重新消费

五、通用代码示例(Spring Boot + Kafka)

电商订单系统的事件驱动实现

5.1 项目结构

src/main/java/com/example/kafka/
├── config/
│   └── KafkaConfig.java          # Kafka配置类
├── event/
│   └── OrderEvent.java           # 事件定义
├── producer/
│   └── OrderEventProducer.java   # 事件生产者
├── consumer/
│   ├── InventoryConsumer.java    # 库存消费者
│   ├── NotificationConsumer.java # 通知消费者
│   └── AnalyticsConsumer.java    # 数据分析消费者
└── controller/
    └── OrderController.java      # 触发入口

src/main/resources/
└── application.yml               # 配置文件

5.2 application.yml

spring:
  kafka:
    # Kafka集群地址
    bootstrap-servers: localhost:9092

    # 生产者配置
    producer:
      # 序列化方式
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      # acks=1 表示 Leader 写入成功即确认
      acks: 1
      # 发送失败重试3次
      retries: 3
      # 批量发送大小:16KB
      batch-size: 16384
      # 等待5ms凑批
      properties:
        linger.ms: 5

    # 消费者配置
    consumer:
      # 反序列化方式
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      # 无offset时从最早开始消费
      auto-offset-reset: earliest
      # 关闭自动提交,手动控制offset
      enable-auto-commit: false
      # 每次最多拉取100条
      max-poll-records: 100

# 自定义Topic名称
app:
  kafka:
    topic:
      order-created: topic-order-created
      order-paid: topic-order-paid
      order-cancelled: topic-order-cancelled

5.3 事件定义

package com.example.kafka.event;

import java.io.Serializable;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.List;

/**
 * 订单事件 - 通用事件结构.
 */
public class OrderEvent implements Serializable {

    /**
     * 事件唯一ID(用于幂等去重).
     */
    private String eventId;

    /**
     * 事件类型:ORDER_CREATED / ORDER_PAID / ORDER_CANCELLED.
     */
    private String eventType;

    /**
     * 事件产生时间.
     */
    private LocalDateTime eventTime;

    /**
     * 事件来源系统.
     */
    private String source;

    /**
     * 业务数据.
     */
    private OrderData data;

    // ----- 内部类:业务数据 -----

    public static class OrderData implements Serializable {
        private String orderId;
        private Integer userId;
        private BigDecimal totalAmount;
        private String status;
        private List<OrderItem> items;

        // getter/setter 省略
        public String getOrderId() { return orderId; }
        public void setOrderId(String orderId) { this.orderId = orderId; }
        public Integer getUserId() { return userId; }
        public void setUserId(Integer userId) { this.userId = userId; }
        public BigDecimal getTotalAmount() { return totalAmount; }
        public void setTotalAmount(BigDecimal totalAmount) { this.totalAmount = totalAmount; }
        public String getStatus() { return status; }
        public void setStatus(String status) { this.status = status; }
        public List<OrderItem> getItems() { return items; }
        public void setItems(List<OrderItem> items) { this.items = items; }
    }

    public static class OrderItem implements Serializable {
        private String skuId;
        private Integer quantity;
        private BigDecimal price;

        // getter/setter 省略
        public String getSkuId() { return skuId; }
        public void setSkuId(String skuId) { this.skuId = skuId; }
        public Integer getQuantity() { return quantity; }
        public void setQuantity(Integer quantity) { this.quantity = quantity; }
        public BigDecimal getPrice() { return price; }
        public void setPrice(BigDecimal price) { this.price = price; }
    }

    // ----- getter/setter -----
    public String getEventId() { return eventId; }
    public void setEventId(String eventId) { this.eventId = eventId; }
    public String getEventType() { return eventType; }
    public void setEventType(String eventType) { this.eventType = eventType; }
    public LocalDateTime getEventTime() { return eventTime; }
    public void setEventTime(LocalDateTime eventTime) { this.eventTime = eventTime; }
    public String getSource() { return source; }
    public void setSource(String source) { this.source = source; }
    public OrderData getData() { return data; }
    public void setData(OrderData data) { this.data = data; }
}

5.4 事件生产者

package com.example.kafka.producer;

import com.alibaba.fastjson.JSON;
import com.example.kafka.event.OrderEvent;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Component;
import org.springframework.util.concurrent.ListenableFuture;
import org.springframework.util.concurrent.ListenableFutureCallback;

import java.time.LocalDateTime;
import java.util.UUID;

/**
 * 订单事件生产者.
 * 负责将订单相关事件发布到Kafka.
 */
@Component
public class OrderEventProducer {

    private static final Logger log = LoggerFactory.getLogger(OrderEventProducer.class);

    private final KafkaTemplate<String, String> kafkaTemplate;

    @Value("${app.kafka.topic.order-created}")
    private String orderCreatedTopic;

    @Value("${app.kafka.topic.order-paid}")
    private String orderPaidTopic;

    @Value("${app.kafka.topic.order-cancelled}")
    private String orderCancelledTopic;

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

    /**
     * 发布订单创建事件.
     *
     * @param orderData 订单数据
     */
    public void publishOrderCreated(OrderEvent.OrderData orderData) {
        OrderEvent event = buildEvent("ORDER_CREATED", orderData);
        sendEvent(orderCreatedTopic, orderData.getOrderId(), event);
    }

    /**
     * 发布订单支付事件.
     *
     * @param orderData 订单数据
     */
    public void publishOrderPaid(OrderEvent.OrderData orderData) {
        OrderEvent event = buildEvent("ORDER_PAID", orderData);
        sendEvent(orderPaidTopic, orderData.getOrderId(), event);
    }

    /**
     * 发布订单取消事件.
     *
     * @param orderData 订单数据
     */
    public void publishOrderCancelled(OrderEvent.OrderData orderData) {
        OrderEvent event = buildEvent("ORDER_CANCELLED", orderData);
        sendEvent(orderCancelledTopic, orderData.getOrderId(), event);
    }

    /**
     * 构建事件对象.
     */
    private OrderEvent buildEvent(String eventType, OrderEvent.OrderData data) {
        OrderEvent event = new OrderEvent();
        event.setEventId(UUID.randomUUID().toString());
        event.setEventType(eventType);
        event.setEventTime(LocalDateTime.now());
        event.setSource("order-service");
        event.setData(data);
        return event;
    }

    /**
     * 发送事件到Kafka.
     * 使用orderId作为key,保证同一订单的事件发送到同一Partition(顺序性).
     *
     * @param topic   目标Topic
     * @param key     消息Key(用于分区路由)
     * @param event   事件对象
     */
    private void sendEvent(String topic, String key, OrderEvent event) {
        String message = JSON.toJSONString(event);
        log.info("发送事件, topic: {}, key: {}, eventType: {}, eventId: {}",
            topic, key, event.getEventType(), event.getEventId());

        // send方法是异步的,返回ListenableFuture
        ListenableFuture<SendResult<String, String>> future =
            kafkaTemplate.send(topic, key, message);

        // 注册回调,处理发送成功和失败
        future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
            @Override
            public void onSuccess(SendResult<String, String> result) {
                log.info("事件发送成功, topic: {}, partition: {}, offset: {}, eventId: {}",
                    result.getRecordMetadata().topic(),
                    result.getRecordMetadata().partition(),
                    result.getRecordMetadata().offset(),
                    event.getEventId());
            }

            @Override
            public void onFailure(Throwable ex) {
                log.error("事件发送失败, topic: {}, eventId: {}, error: {}",
                    topic, event.getEventId(), ex.getMessage(), ex);
                // 这里可以做补偿:写入本地表、重试队列等
            }
        });
    }
}

5.5 事件消费者

package com.example.kafka.consumer;

import com.alibaba.fastjson.JSON;
import com.example.kafka.event.OrderEvent;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;

/**
 * 库存消费者.
 * 监听订单创建事件,执行库存预占.
 * 监听订单取消事件,执行库存释放.
 */
@Component
public class InventoryConsumer {

    private static final Logger log = LoggerFactory.getLogger(InventoryConsumer.class);

    /**
     * 消费订单创建事件 - 执行库存预占.
     *
     * topics: 监听的Topic
     * groupId: 消费者组(同组内负载均衡,不同组各自消费全量)
     * containerFactory: 使用手动提交offset的容器工厂
     */
    @KafkaListener(
        topics = "${app.kafka.topic.order-created}",
        groupId = "inventory-service"
    )
    public void onOrderCreated(ConsumerRecord<String, String> record, Acknowledgment ack) {
        try {
            log.info("库存服务收到订单创建事件, partition: {}, offset: {}, key: {}",
                record.partition(), record.offset(), record.key());

            // 反序列化事件
            OrderEvent event = JSON.parseObject(record.value(), OrderEvent.class);

            // 幂等校验:根据eventId判断是否已处理过
            if (isEventProcessed(event.getEventId())) {
                log.info("事件已处理,跳过, eventId: {}", event.getEventId());
                ack.acknowledge(); // 提交offset
                return;
            }

            // 执行库存预占
            for (OrderEvent.OrderItem item : event.getData().getItems()) {
                reserveStock(item.getSkuId(), item.getQuantity());
            }

            // 记录已处理的事件ID(幂等表)
            markEventProcessed(event.getEventId());

            // 手动提交offset,表示消息已成功处理
            ack.acknowledge();
            log.info("库存预占成功, orderId: {}", event.getData().getOrderId());

        } catch (Exception e) {
            log.error("库存预占失败, record: {}", record.value(), e);
            // 不提交offset,消息会在下次poll时重新消费(At Least Once语义)
            // 生产中可以:重试N次后发送到死信队列
        }
    }

    /**
     * 消费订单取消事件 - 执行库存释放.
     */
    @KafkaListener(
        topics = "${app.kafka.topic.order-cancelled}",
        groupId = "inventory-service"
    )
    public void onOrderCancelled(ConsumerRecord<String, String> record, Acknowledgment ack) {
        try {
            OrderEvent event = JSON.parseObject(record.value(), OrderEvent.class);
            log.info("收到订单取消事件, orderId: {}", event.getData().getOrderId());

            // 释放库存
            for (OrderEvent.OrderItem item : event.getData().getItems()) {
                releaseStock(item.getSkuId(), item.getQuantity());
            }

            ack.acknowledge();
            log.info("库存释放成功, orderId: {}", event.getData().getOrderId());
        } catch (Exception e) {
            log.error("库存释放失败", e);
        }
    }

    // ----- 模拟业务方法 -----

    private boolean isEventProcessed(String eventId) {
        // 实际:查询幂等表 SELECT 1 FROM event_log WHERE event_id = ?
        return false;
    }

    private void markEventProcessed(String eventId) {
        // 实际:INSERT INTO event_log(event_id, process_time) VALUES(?, NOW())
    }

    private void reserveStock(String skuId, Integer quantity) {
        log.info("预占库存: skuId={}, quantity={}", skuId, quantity);
        // 实际:UPDATE stock SET reserved = reserved + ? WHERE sku_id = ?
    }

    private void releaseStock(String skuId, Integer quantity) {
        log.info("释放库存: skuId={}, quantity={}", skuId, quantity);
        // 实际:UPDATE stock SET reserved = reserved - ? WHERE sku_id = ?
    }
}

package com.example.kafka.consumer;

import com.alibaba.fastjson.JSON;
import com.example.kafka.event.OrderEvent;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;

/**
 * 通知消费者.
 * 监听订单事件,发送短信/推送通知给用户.
 * 独立的消费者组,与库存服务互不影响.
 */
@Component
public class NotificationConsumer {

    private static final Logger log = LoggerFactory.getLogger(NotificationConsumer.class);

    @KafkaListener(
        topics = "${app.kafka.topic.order-created}",
        groupId = "notification-service"  // 不同的groupId,独立消费全量消息
    )
    public void onOrderCreated(ConsumerRecord<String, String> record, Acknowledgment ack) {
        try {
            OrderEvent event = JSON.parseObject(record.value(), OrderEvent.class);
            Integer userId = event.getData().getUserId();
            String orderId = event.getData().getOrderId();

            // 发送下单成功通知
            sendSms(userId, "您的订单 " + orderId +
                        // 发送下单成功通知
            sendSms(userId, "您的订单 " + orderId + " 已创建成功,请尽快支付。");
            sendAppPush(userId, "下单成功", "订单" + orderId + "等待支付");

            ack.acknowledge();
            log.info("通知发送成功, userId: {}, orderId: {}", userId, orderId);
        } catch (Exception e) {
            log.error("通知发送失败", e);
            // 通知失败不需要重试太多次,可以直接ack跳过
            ack.acknowledge();
        }
    }

    @KafkaListener(
        topics = "${app.kafka.topic.order-paid}",
        groupId = "notification-service"
    )
    public void onOrderPaid(ConsumerRecord<String, String> record, Acknowledgment ack) {
        try {
            OrderEvent event = JSON.parseObject(record.value(), OrderEvent.class);
            Integer userId = event.getData().getUserId();
            String orderId = event.getData().getOrderId();

            sendSms(userId, "您的订单 " + orderId + " 支付成功,正在为您安排发货。");

            ack.acknowledge();
        } catch (Exception e) {
            log.error("支付通知发送失败", e);
            ack.acknowledge();
        }
    }

    // ----- 模拟方法 -----

    private void sendSms(Integer userId, String content) {
        log.info("发送短信, userId: {}, content: {}", userId, content);
    }

    private void sendAppPush(Integer userId, String title, String body) {
        log.info("发送APP推送, userId: {}, title: {}", userId, title);
    }
}
package com.example.kafka.consumer;

import com.alibaba.fastjson.JSON;
import com.example.kafka.event.OrderEvent;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;

/**
 * 数据分析消费者.
 * 监听所有订单事件,做实时统计和数据分析.
 * 独立消费者组,失败不影响任何业务.
 */
@Component
public class AnalyticsConsumer {

    private static final Logger log = LoggerFactory.getLogger(AnalyticsConsumer.class);

    @KafkaListener(
        topics = {
            "${app.kafka.topic.order-created}",
            "${app.kafka.topic.order-paid}",
            "${app.kafka.topic.order-cancelled}"
        },
        groupId = "analytics-service"
    )
    public void onOrderEvent(ConsumerRecord<String, String> record, Acknowledgment ack) {
        try {
            OrderEvent event = JSON.parseObject(record.value(), OrderEvent.class);

            switch (event.getEventType()) {
                case "ORDER_CREATED":
                    incrementCounter("order.created.count");
                    recordAmount("order.created.amount", event.getData().getTotalAmount());
                    break;
                case "ORDER_PAID":
                    incrementCounter("order.paid.count");
                    recordAmount("order.paid.amount", event.getData().getTotalAmount());
                    break;
                case "ORDER_CANCELLED":
                    incrementCounter("order.cancelled.count");
                    break;
                default:
                    log.warn("未知事件类型: {}", event.getEventType());
            }

            ack.acknowledge();
        } catch (Exception e) {
            log.error("数据分析处理异常", e);
            ack.acknowledge(); // 分析服务容忍数据丢失,直接跳过
        }
    }

    // ----- 模拟统计方法 -----

    private void incrementCounter(String metric) {
        log.info("统计指标+1: {}", metric);
    }

    private void recordAmount(String metric, Object amount) {
        log.info("记录金额: {} = {}", metric, amount);
    }
}

5.6 Kafka 配置类(手动提交 offset)

package com.example.kafka.config;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.ContainerProperties;

import java.util.HashMap;
import java.util.Map;

/**
 * Kafka消费者配置.
 * 配置手动提交offset模式.
 */
@Configuration
public class KafkaConfig {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    /**
     * 消费者工厂配置.
     */
    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        // 关闭自动提交
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);
        return new DefaultKafkaConsumerFactory<>(props);
    }

    /**
     * 监听容器工厂 - 手动提交模式.
     */
    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        // 设置手动提交offset模式
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
        // 并发消费者数量(对应Partition数量)
        factory.setConcurrency(3);
        return factory;
    }
}

5.7 触发入口(Controller)

package com.example.kafka.controller;

import com.example.kafka.event.OrderEvent;
import com.example.kafka.producer.OrderEventProducer;
import org.springframework.web.bind.annotation.*;

import java.math.BigDecimal;
import java.util.Arrays;

/**
 * 订单接口 - 演示事件发布.
 */
@RestController
@RequestMapping("/api/orders")
public class OrderController {

    private final OrderEventProducer eventProducer;

    public OrderController(OrderEventProducer eventProducer) {
        this.eventProducer = eventProducer;
    }

    /**
     * 创建订单.
     * 业务处理完成后,发布"订单创建"事件.
     */
    @PostMapping
    public String createOrder() {
        // 1. 业务逻辑:创建订单、写数据库
        String orderId = "ORD-" + System.currentTimeMillis();

        // 2. 组装事件数据
        OrderEvent.OrderData orderData = new OrderEvent.OrderData();
        orderData.setOrderId(orderId);
        orderData.setUserId(10001);
        orderData.setTotalAmount(new BigDecimal("299.00"));
        orderData.setStatus("CREATED");

        OrderEvent.OrderItem item = new OrderEvent.OrderItem();
        item.setSkuId("SKU-001");
        item.setQuantity(2);
        item.setPrice(new BigDecimal("149.50"));
        orderData.setItems(Arrays.asList(item));

        // 3. 发布事件(异步,不阻塞主流程)
        eventProducer.publishOrderCreated(orderData);

        // 4. 立即返回用户
        return "订单创建成功: " + orderId;
    }

    /**
     * 取消订单.
     */
    @PostMapping("/{orderId}/cancel")
    public String cancelOrder(@PathVariable String orderId) {
        // 1. 业务逻辑:更新订单状态为已取消

        // 2. 发布取消事件
        OrderEvent.OrderData orderData = new OrderEvent.OrderData();
        orderData.setOrderId(orderId);
        orderData.setUserId(10001);

        OrderEvent.OrderItem item = new OrderEvent.OrderItem();
        item.setSkuId("SKU-001");
        item.setQuantity(2);
        orderData.setItems(Arrays.asList(item));

        eventProducer.publishOrderCancelled(orderData);

        return "订单已取消: " + orderId;
    }
}

六、Kafka 关键流程操作

6.1 消息发送流程(Producer 端)

应用调用 send()
      │
      ▼
┌─────────────────────┐
│ 拦截器(Interceptor)│  可选,做日志、监控等
└──────────┬──────────┘
           │
           ▼
┌─────────────────────┐
│ 序列化(Serializer) │  对象 → 字节数组
└──────────┬──────────┘
           │
           ▼
┌─────────────────────┐
│ 分区器(Partitioner)│  决定发到哪个Partition
│                     │  · 有key:hash(key) % partition数
│                     │  · 无key:轮询
└──────────┬──────────┘
           │
           ▼
┌─────────────────────┐
│ 消息累加器           │  按 Partition 分组缓存消息
│ (RecordAccumulator) │  攒批量(batch.size 或 linger.ms)
└──────────┬──────────┘
           │
           ▼
┌─────────────────────┐
│ Sender 线程          │  后台线程,批量发送到 Broker
│                     │  · 建立TCP连接
│                     │  · 按Broker分组发送
└──────────┬──────────┘
           │
           ▼
┌─────────────────────┐
│ Broker 接收         │  写入 Partition 日志
│                     │  分配 Offset
│                     │  返回 ack
└─────────────────────┘

6.2 消息消费流程(Consumer 端)

Consumer 启动
      │
      ▼
┌─────────────────────────┐
│ 加入 Consumer Group      │
│ (JoinGroup 协议)         │
│ · Coordinator 分配分区   │
│ · Rebalance(重平衡)    │
└──────────┬──────────────┘
           │
           ▼
┌─────────────────────────┐
│ 获取初始 Offset          │
│ · 有提交记录:从上次位置 │
│ · 无记录:根据策略        │
│   earliest / latest      │
└──────────┬──────────────┘
           │
           ▼
  ┌────────────────┐
  │  poll() 循环    │ ◄───────────────────┐
  │                 │                      │
  │  发送fetch请求  │                      │
  │  ↓              │                      │
  │  拉取一批消息   │                      │
  │  ↓              │                      │
  │  反序列化       │                      │
  │  ↓              │                      │
  │  业务处理       │                      │
  │  ↓              │                      │
  │  提交 offset    │ ─────────────────────┘
  └────────────────┘

6.3 Rebalance(重平衡)

当以下事件发生时,Kafka 会重新分配 Partition 给 Consumer:

触发条件:
· Consumer 加入或离开 Group
· Consumer 心跳超时(被认为宕机)
· Topic 的 Partition 数量变化
· Consumer 调用 subscribe() 订阅新 Topic

重平衡过程(会导致短暂消费暂停):
1. 所有 Consumer 停止消费
2. Group Coordinator 收集所有 Consumer 的信息
3. Leader Consumer 执行分区分配策略
4. 每个 Consumer 获得新的 Partition 分配方案
5. 恢复消费

七、关键生产问题与解决方案

7.1 消息丢失

丢失位置 原因 解决方案
Producer 端 acks=0,不等确认 acks=all + retries > 0
Broker 端 Leader 宕机且副本未同步完 min.insync.replicas=2
Consumer 端 自动提交 offset 后处理失败 手动提交 offset(处理成功再提交)

7.2 消息重复

重复位置 原因 解决方案
Producer 端 网络超时重发 开启幂等(enable.idempotence=true)
Consumer 端 处理成功但 offset 提交失败 消费端幂等(唯一ID + 去重表)

7.3 消息顺序

场景 方案
全局有序 单 Partition(牺牲并行度)
局部有序 相同 key 的消息发到同一 Partition
不需要顺序 多 Partition + 多 Consumer 并行消费

7.4 消费积压

手段 说明
增加 Partition 数 提高并行度上限
增加 Consumer 数 同组内消费者数 ≤ Partition 数
提高 max.poll.records 每次多拉一些消息
批量处理 拉取后批量写入 DB,减少 IO 次数
临时扩容 高峰期动态增加消费者实例

八、Kafka 在事件驱动中的最佳实践

8.1 事件设计原则

一个好的事件应该包含:
┌─────────────────────────────────────┐
│ 事件信封(Envelope)                │
│ ├── eventId:全局唯一标识(幂等用) │
│ ├── eventType:事件类型             │
│ ├── eventTime:事件发生时间         │
│ ├── source:来源系统                │
│ └── data:业务数据                  │
│     ├── 包含消费者需要的所有信息    │
│     └── 避免消费者再反查生产者      │
└─────────────────────────────────────┘

8.2 Topic 命名规范

推荐格式:{领域}.{实体}.{动作}
示例:
· order.created       订单创建
· payment.completed   支付完成
· stock.reserved      库存预占
· user.registered     用户注册

或者:{前缀}-{环境}-{业务}
示例:
· topic-order-events-dev
· topic-stock-change-prod

8.3 Consumer Group 规划

原则:一个微服务对应一个 Consumer Group

Topic: order-events
├── group: inventory-service    → 库存服务(独立消费全量)
├── group: notification-service → 通知服务(独立消费全量)
├── group: analytics-service    → 分析服务(独立消费全量)
└── group: audit-service        → 审计服务(独立消费全量)

每个服务独立消费、独立进度、互不影响
某个服务消费失败,不影响其他服务

九、总结对比

维度 同步直接调用 Kafka 事件驱动
耦合度 高(A 必须知道 B 的接口) 低(A 只管发事件)
响应时间 所有下游耗时叠加 仅主流程耗时
可用性 任何下游挂了都失败 下游挂了不影响上游
扩展性 新增下游要改上游代码 新增消费者即可
数据一致性 强一致(事务) 最终一致(需要幂等设计)
调试难度 简单(同步链路清晰) 较复杂(异步链路需要追踪)
适用场景 强一致要求、简单系统 高吞吐、微服务、解耦需求强

一句话总结:Kafka 就是系统间的「高速公路」——生产者把货(消息)放上去就走,消费者按自己的速度从上面取货。公路有多车道(Partition)保证高吞吐,有收费站记录(Offset)保证不丢货,有备份路线(Replica)保证不断路。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐