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=1Leader 写入即确认,平衡性能和可靠性
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编程工具,助力开发者即刻编程。

更多推荐