欢迎光临
我们一直在努力

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

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


一、为什么需要 Kafka

先从一个实际问题说起:

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

  • 扣减库存
  • 通知仓库发货
  • 发送短信
  • 更新积分
  • 同步数据到第三方系统
  • 同步调用的问题:

    用户下单 → 扣库存(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: topicordercreated
    order-paid: topicorderpaid
    order-cancelled: topicordercancelled

    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)保证不断路。

    赞(0)
    未经允许不得转载:171主机测评 » Kafka 异步消息推送与事件驱动架构详解
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址