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: 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)保证不断路。
更多推荐




所有评论(0)