kafka基础和应用
前言
- 会随着经验不断更新这篇博客,关于kafka的运用和实践。
1.什么是kafka
- Kafka 是一个分布式、可持久化、可重放的日志流系统,用于解耦系统间的数据流转。
- kafka不是传统意义上的消息队列,传统的消息队列在消费完数据后会直接删除,而kafka在消费后不会被删除.
- kafka是一个流平台,流平台不仅仅是数据流转,而是基于时间顺序追加、可持久化、可回放的数据流模型。
2.它解决了什么问题
- 它解决的核心问题就是
专门负责不同系统之间的传递数据,支持有顺序、可持久化、可回放
1.异步解耦
- 在前面秒杀优惠券场景下,我们通常的做法是通过lua脚本先去查询并扣除库存,然后将消息发送至Stream消息队列中,再由后台多开的线程进行读取消息,进行异步处理,进行落库,而kafka就可以替换redis中的stream消息队列,更加工程化,保证redis单一职责的架构模式
2.削峰填谷
- kafka相当于一个超大容量的缓冲池,能够将流量削平,慢慢消费直到落库
3.持久化
- kafka可以支持数据写磁盘,支持副本机制,宕机可恢复
4.分布式系统间的通信
- kafka适用于分布式场景下,我们之前用的阻塞队是放在jvm里面的,倘若部署在不同的机器上面,jvm都会不同,所以无法进行不同机器之间的交流,而kafka是一个独立的服务器,所有服务器都能访问
5.大规模日志和数据流处理
- 具体遇到我们再吸取经验,先大致了解一下

3.基本概念

- offset 是 Kafka 每条消息在分区里的“编号”。
kafka中的结构
消息结构:
| topic | partition | offset | key | value |
消费状态结构:
| groupId | topic | partition | currentOffset |
partition如果不指定 是根据key算出来的

分区里面的结构

- 发送–消费流程
Producer 通过 Topic 发送消息,Kafka 根据 Key 选择 Partition,然后把消息顺序写入 Partition 日志并生成 Offset。
Consumer 先订阅 Topic,Kafka 再把 Topic 的 Partition 分配给 Consumer,然后 Consumer 按 Offset 顺序消费 Partition 中的消息。
- 消费者组的概念
同一消费者组内:
每个分区只能有一个消费者消费
保证消息只被消费一次(竞争消费)
不同消费者组之间:
每个组都会消费完整数据
实现消息广播
通过增加消费者数量,可以提高系统吞吐能力
- broker是啥
Kafka Broker 是 Kafka 集群中的一个节点,本质就是一台 Kafka 服务器。一个 Kafka 集群通常由多个 Broker 组成。Broker 的主要职责是存储消息数据、接收生产者发送的消息以及提供给消费者读取。在 Kafka 中,Topic 会被划分为多个 Partition,而每个 Partition 会分布在不同的 Broker 上,因此 Broker 本质上是 Partition 的承载节点。同时 Kafka 还会为 Partition 设置副本机制,不同 Broker 之间会进行数据同步(ISR),从而保证高可用和数据可靠性。
4.实现发送消息和接收消息的demo
- 1.如果你是windows系统,下载好docker desktop,选择AMD64,确保能够启动docker服务,如果是linux,直接启动docker就行了,下载kafka

- 2.书写docker-composer.yml
# 表示这是个容器
services:
# zookeeper服务
zookeeper:
# 镜像选择 如果本地没有会自动pull
image: confluentinc/cp-zookeeper:7.6.0
# 环境变量
environment:
# zookeeper对外提供的客户端连接的端口是 2181
ZOOKEEPER_CLIENT_PORT: 2181
# zookeeper的心跳时间单位是 2000 ms -> 2s
ZOOKEEPER_TICK_TIME: 2000
# 宿主机的2181端口 映射到 容器的 2181端口
ports:
- "2181:2181"
# kafka服务
kafka:
# 镜像选择
image: confluentinc/cp-kafka:7.6.0
# 启动顺序 在zookeeper之后 不保证zookeeper启动成功
depends_on:
- zookeeper
# 将宿主机的9092端口 映射到 容器的 9092端口
ports:
- "9092:9092"
# 环境变量
environment:
# kafka的节点brokerId
KAFKA_BROKER_ID: 1
# kafka要连接zookeeper的地址 zookeeper 是 compose 的 service 名(容器网络 DNS 可解析)
KAFKA_ZOOKEEPER_CONNECT: "zookeeper:2181"
# 两套监听:宿主机 + 容器内部
# 因为 Kafka 的客户端既可能来自宿主机,也可能来自容器内部网络,而不同网络环境下访问地址不同,
# 所以必须配置多个监听器和对应的广播地址。
# kafka在容器里的所有网卡都监听9092端口 在容器内部监听29092端口
KAFKA_LISTENERS: "PLAINTEXT://0.0.0.0:9092,PLAINTEXT_INTERNAL://0.0.0.0:29092"
# kafka返回给客户端(对外)的节点地址是advised里面的地址在 本机9092端口 容器里的kafka服务的29092端口使用
KAFKA_ADVERTISED_LISTENERS: "PLAINTEXT://localhost:9092,PLAINTEXT_INTERNAL://kafka:29092"
# kafka的监听器使用明文协议 容器内部也使用明文协议
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: "PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT"
# 集群内部的节点名字 解决数据同步通道问题
KAFKA_INTER_BROKER_LISTENER_NAME: "PLAINTEXT_INTERNAL"
# 消费位移的副本数量 单机必须是 1
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
# kafka事务相关内部主题的副本数
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
# ISR = “正在同步的副本集合”(你可以先粗略理解为“可用副本”) 单机只有一个
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
# consumer group 第一次 rebalance 时的等待时间 ms 重新分配分区的时间间隔
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
- 3.cd进docker-composer目录,docker composer up -d启动
KafkaDemoConsumer
package cn.fly.kafkademo.consumer;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Component
public class KafkaDemoConsumer {
@KafkaListener(topics = "demo-topic", groupId = "demo-group")
public void onMessage(String msg) {
System.out.println("我监听到了:topic为:demo-topic groupId为:demo-group的kafka消息");
System.out.println("✅ consumed: " + msg);
}
}
KafkaDemoController
package cn.fly.kafkademo.controller;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.web.bind.annotation.*;
import lombok.*;
@RestController
@RequiredArgsConstructor
@RequestMapping("/kafka")
public class KafkaDemoController {
private final KafkaTemplate<String, String> kafkaTemplate;
@PostMapping("/send")
public String send(@RequestParam String msg) {
System.out.println(msg);
kafkaTemplate.send("demo-topic", msg);
return "sent: " + msg;
}
}
application.yml
spring:
kafka:
# kafka集群的入口
bootstrap-servers: localhost:9092
# 当前消费者属于哪个消费者组
consumer:
group-id: demo-group
# 当offset不存在时怎样读消息 earliest最早 latest最新 none报错
auto-offset-reset: earliest
producer:
# 生产者发送消息后多少节点确认才算成功
acks: all
pom.xml
<dependencies>
<!-- kafka依赖-->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<!-- Lombok -->
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<!-- Web(做接口测试Kafka用) -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>
5.常见问题
1.你的maven是否配置正确
- 如果你的maven配置正确,还是启动不了,还是有依赖没导入,打开maven的生命周期,clean一遍,再compile一遍,重新加载maven,就可以了
2.为什么要双监听
- listeners 决定 Kafka 实际监听的端口;advertised.listeners 决定 Kafka 返回给客户端的可达地址。单机时两者常可一样,但在 Docker/内外网/NAT 场景下必须区分,否则客户端会拿到不可达的 broker 地址导致连接失败。比如 0.0.0.0:9092 表示容器里所有网卡都监听 9092。
3.什么是网卡
-
有一个误区,网卡和ip地址难道不一样吗?
-
不一样,网卡是一个网络接口,一个连通网络的实际通道,而ip地址就是这个通道的地址,网卡就是实际的房子,ip就是门牌号,同一个网卡可以有多个ip,而每个ip只能对应一个网卡
-
比如:0.0.0.0:9092 表示所有网卡的所有ip都可以监听9092端口,192.168.1.10:9092 表示只有这个ip对应的那个网卡才能监听9092端口
4.那既然listener已经能监听到kafka了 那为啥还要告诉客户端用哪个地址来连接我?
- listeners 决定“我在哪接电话”
advertised.listeners 决定“我告诉别人我的电话是多少”,在分布式场景下,我的kafka假设我部署在公司私网地址里面,那客户端来访问kafka监听的端口,这个时候advertisedListener返回客户端你应该访问的地址,如果这个地址还是公司私网,那客户端根本连接不了,这说明,advertisedListener解决了kafka监听的地址,可能和kafka告诉客户端应该访问的地址不一致,导致客户端无法访问的问题

5.但是假设我的kafka就只部署在一台机器上呢?
- 如果kafka只部署在一台机器上,那么就可以一样,监听地址 = 客户端可达地址
6.埋点
什么是埋点
- 埋点就是在关键业务记录业务节点行为数据
它解决了什么问题
- 它能够统计用户的行为,进行数据分析,监控系统异常
它的分类
- 它分为业务埋点、日志埋点、链路埋点
- 业务埋点记录用户的行为,比如下单、点击、支付
- 日志埋点记录异常、性能、耗时
- 链路埋点记录调用链
为什么要用kafka来做埋点
- 因为我们要记录用户的行为和数据分析,如果我们每时每刻都直接写库,数据库的QPS(每秒请求量)就会巨多,而kafka时高吞吐的顺序日志系统,
kafka把埋点行为从同步操作变成了异步数据流
什么是异步数据流
- 异步数据流 = 数据产生后,不立刻处理,而是先进入“流”,由其他系统稍后处理。
- 数据流是一条源源不断产生的数据河
7.埋点的实际应用demo
- 下单埋点操作,打开自己的kafkaDemo 项目
生产者在传入消息的时候没有指定分区,为什么消费者能够从该分区里读到消息
- 因为生产者只负责生产消息,具体消费者怎么读,怎么分区读,与生产者无关,并且每一个不同的分区,都有自己独立的offset,偏移量,
kafka的本质只是把消息写入日志,仅此而已,发给谁,与kafka无关
topic 3 个 partition
groupA 有 5 个消费者
groupB 有 1 个消费者
一条消息会被消费几次?
- groupA 中只有 3 个消费者会被分配分区(2 个空闲),但对同一条消息只会有其中 1 个消费者消费;groupB 也会消费该消息一次,所以总共消费 2 次。
8.实战演练
- 使用 Redis Lua 脚本实现库存原子扣减与一人一单校验,结合 Kafka 实现异步削峰与订单落库,通过数据库唯一索引保证幂等消费
- 打开自己的kafkaDemo,Seckill相关
8.1 Result类
在这个Result结构中的问题:
1.为什么要实现Serializable接口
- 因为主要是为了
支持对象序列化,虽然在springwebjson返回场景中不是必须,但作为通用DTO通常会实现。
2.@Data注解是不是包含了@AllArgsConstructor 和@NoArgsConstructor
- @Data注解
不包含全参构造器和无参构造器,它只包含getter、setter、toString、equals、hashCode和@RequiredArgsConstructor
3.方法中为什么要返回 Result而不是Result
- 因为
静态方法无法直接使用类级泛型T,所以必须在方法前重新声明泛型T
4.为什么静态方法不能使用类级泛型T?
- 因为
类的泛型T是在创建实例时确定的,而静态方法在类加载时就已经确定了,它不依赖实例,因此不能访问类级泛型
8.2 controller
1.总是忘记接收对象时需要@RequestBody进行接收
2.传入对象的时候注意验空,前端也会看到,@Valid和@NotNull
3.设置了@Valid和@NotNull需要用全局异常来捕获
8.3 exception
1.什么是@RestControllerAdvice
- @RestControllerAdvice 是一个
全局异常增强组件,它会拦截所有 Controller 抛出的异常,并统一处理成 JSON 响应,它等于@ControllerAdvice + @ResponseBody
2.@ExceptionHandler(MethodArgumentNotValidException.class)是啥
- 当出现指定异常MethodArgumentNotValidException时,执行这个方法。
3.Result<?>为什么不是Result<String>
- 使用 Result<?> 是因为异常处理器是通用的,不限定 data 的具体类型,使用通配符泛型可以提高通用性。
4.那为啥是?而不是T
- 因为这里没有泛型类型参数 T 的定义,不能凭空使用 T。
/**
* 全局异常捕获类
*/
@RestControllerAdvice
public class GlobalExceptionHandler {
@ExceptionHandler(MethodArgumentNotValidException.class)
public Result<?> handleValidException(MethodArgumentNotValidException e) {
String msg = Objects.requireNonNull(e.getBindingResult()
.getFieldError())
.getDefaultMessage();
return Result.fail(msg);
}
}
?和T的区别:

请求执行流程:

8.4 建表
# 建库
create database if not exists seckill;
# 使用这个库
use seckill;
# 建表
create table if not exists seckillOrder (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
order_id BIGINT NOT NULL,
voucher_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_order_id (order_id),
UNIQUE KEY uk_voucher_user (voucher_id, user_id)
) ENGINE=INNODB;
1.PRIMARY_KEY 在INNODB引擎中的主键是
聚簇索引
2.什么是
聚簇索引
- 聚簇索引是
数据本身就存储在索引里的一种索引

3.INCREMENT的作用
- 主键自增,按顺序插入,减少B+树的重排,对B+树友好
4.MyISAM 是 MySQL 早期的存储引擎
- 它不支持事务、行锁、外键、查询快但是写并发差,在秒杀场景下需要写,事务和行锁进行控制
5 定义索引:UNIQUE KEY 索引名 (字段名)
UNIQUE KEY uk_order_id (order_id),
- 它等价于
CREATE UNIQUE INDEX uk_order_id ON seckillOrder(order_id);
6.聚簇索引解决了什么问题
- 加快了查询的效率

7.到底什么是索引,它解决了什么问题
- 索引是
帮助数据表快速查找数据的一种数据结构,它在Innodb引擎里面是B+树结构,它能够将数据插入B+树中,查找时间复杂度从O(n)变为O(logn)
暂时到这为止,数据库的知识在另一个笔记里面
8.5 lua脚本编写
-- KEYS[1] stockKey: seckill:stock:{voucherId}
-- KEYS[2] userSetKey: seckill:order:users:{voucherId}
-- ARGV[1] voucherId
-- ARGV[2] userId
-- ARGV[3] orderId (generated by Java)
-- 0 成功 1 库存不足 2 重复下单
local stockKey = KEYS[1]
local userSetKey = KEYS[2]
local userId = ARGV[2]
local orderId = ARGV[3]
-- 一人一单
-- SISMEMBER userSetKey userId = 1 ?
if redis.call('SISMEMBER', userSetKey, userId) == 1 then
-- 返回json
return 2
end
-- 库存校验 防止超卖
-- string num = GET stockKey -> int num = tonumber(num)
local stock = tonumber(redis.call('GET', stockKey))
if(stock == nil) or (stock <= 0) then
return 1
end
-- 扣库存 记录用户
-- DECR stockKey 让stockKey --
redis.call('DECR', stockKey)
-- SADD userSetKey userId 在userSet里面添加userId
redis.call('SADD', userSetKey, userId)
return 0
1.redis Key的书写规范
业务:模块:含义:变量
seckill:order:users:1
2.KEYS是什么
- 是redis的key键
3.ARGV是什么
- 是传入lua脚本的参数,下标从1开始
8.6 maven
- maven的地址和本地仓库的结构:


<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<!-- 关键:Spring Boot parent 提供 dependencyManagement(版本托管) -->
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.3.2</version>
<relativePath/>
</parent>
<groupId>com.example</groupId>
<artifactId>kafkaDemo</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<java.version>17</java.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
/* 注意这个kafka依赖不是starter 已经拉不到了 我卡了1个小时 服了*/
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-validation</artifactId>
</dependency>
<!-- MapStruct -->
<dependency>
<groupId>org.mapstruct</groupId>
<artifactId>mapstruct</artifactId>
<version>1.5.5.Final</version>
</dependency>
<!-- 编译期代码生成 -->
<dependency>
<groupId>org.mapstruct</groupId>
<artifactId>mapstruct-processor</artifactId>
<version>1.5.5.Final</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
8.7 consumer
第1版
@Slf4j
@Component
@RequiredArgsConstructor
public class SeckillConsumer {
private final ObjectMapper objectMapper = new ObjectMapper();
private final SeckillOrderMapper seckillOrderMapper;
@KafkaListener(topics = KafkaConstants.SECKILL_TOPIC)
public void seckill(ConsumerRecord<String, String> record, Acknowledgment ack) throws Exception {
String json = record.value();
// 解析消息
OrderMsg msg = objectMapper.readValue(json, OrderMsg.class);
LambdaQueryWrapper<SeckillOrder> queryWrapper = Wrappers.lambdaQuery();
queryWrapper.eq(SeckillOrder::getOrderId, msg.orderId)
.eq(SeckillOrder::getVoucherId, msg.voucherId)
.eq(SeckillOrder::getUserId, msg.userId);
Long count = seckillOrderMapper.selectCount(queryWrapper);
// 如果已经存在该订单就不需要插入了
if (count > 0) {
log.info("seckill存在重复的订单, orderId: {}", msg.orderId);
// 手动提交offset让kafka不要再重复发消息了
ack.acknowledge();
return;
}
// 如果不存在那就插入
StoreOrderBasicVo storeOrderBasicVo = new StoreOrderBasicVo();
storeOrderBasicVo.setOrderId(msg.orderId);
storeOrderBasicVo.setVoucherId(msg.voucherId);
storeOrderBasicVo.setUserId(msg.userId);
seckillOrderMapper.insert(StoreOrderConverter.INSTANCE.basicVoToPo(storeOrderBasicVo));
// 成功才 ack
ack.acknowledge();
}
// 你消息里就是 orderId/voucherId/userId,这里做个内部类就够了
public static class OrderMsg {
public long orderId;
public long voucherId;
public long userId;
}
}
第一版的问题:
- 1.先判断再插入就能保证幂等性吗?万一两个线程同时进入条件判断,它们查到的都是0,同时插入,就会重复下单,这和我们之前为什么要用lua脚本一摸一样,只不过这里只能用数据库层面的约束了
- 2.那为啥不用lua脚本保证原子性?因为lua是redis层面的,这里是consumer马上就要落库了,不适合做幂等,如果用了,反而会再多一次io操作,等待lua脚本返回结果,而且lua脚本到处都在用,耦合性高,不适合当前场景
第2版
@Slf4j
@Component
@RequiredArgsConstructor
public class SeckillConsumer {
private final ObjectMapper objectMapper = new ObjectMapper();
private final SeckillOrderMapper seckillOrderMapper;
@KafkaListener(topics = KafkaConstants.SECKILL_TOPIC)
public void seckill(ConsumerRecord<String, String> record, Acknowledgment ack) throws Exception {
String json = record.value();
// 解析消息
OrderMsg msg = objectMapper.readValue(json, OrderMsg.class);
// 直接利用数据库唯一键约束进行幂等
StoreOrderBasicVo storeOrderBasicVo = new StoreOrderBasicVo();
storeOrderBasicVo.setOrderId(msg.orderId);
storeOrderBasicVo.setVoucherId(msg.voucherId);
storeOrderBasicVo.setUserId(msg.userId);
try{
seckillOrderMapper.insert(StoreOrderConverter.INSTANCE.basicVoToPo(storeOrderBasicVo));
// 注意捕获的是唯一键约束 我们的唯一键索引是建立在orderId上面
} catch (DuplicateKeyException e) {
log.info("seckill存在重复消息: orderId: {}", storeOrderBasicVo.getOrderId());
// 防止kafka一直重试
ack.acknowledge();
}
// 成功才 ack
ack.acknowledge();
}
// 你消息里就是 orderId/voucherId/userId,这里做个内部类就够了
public static class OrderMsg {
public long orderId;
public long voucherId;
public long userId;
}
}
- 1.注意捕获的是唯一键约束异常,因为其他错误如果也被捕获了的话直接ack了,会导致这条消息的丢失!
- 2.那么索引就是唯一键吗?mysql里面有几种索引?

第3版
@Slf4j
@Component
@RequiredArgsConstructor
public class SeckillConsumer {
private final SeckillOrderMapper seckillOrderMapper;
private final ObjectMapper objectMapper;
@KafkaListener(topics = KafkaConstants.SECKILL_TOPIC)
public void seckill(ConsumerRecord<String, String> record, Acknowledgment ack) throws Exception {
String json = record.value();
// 解析消息
EventMessage<SeckillOrderCreatedEvent> message = objectMapper.readValue(json,
new TypeReference<>() {
});
SeckillOrderCreatedEvent msg = message.getData();
// 直接利用数据库唯一键约束进行幂等
StoreOrderBasicVo storeOrderBasicVo = new StoreOrderBasicVo();
storeOrderBasicVo.setOrderId(msg.getOrderId());
storeOrderBasicVo.setVoucherId(msg.getVoucherId());
storeOrderBasicVo.setUserId(msg.getUserId());
try{
seckillOrderMapper.insert(StoreOrderConverter.INSTANCE.basicVoToPo(storeOrderBasicVo));
// 注意捕获的是唯一键约束 我们的唯一键索引是建立在orderId上面
} catch (DuplicateKeyException e) {
log.info("seckill存在重复消息: orderId: {}", storeOrderBasicVo.getOrderId());
// 防止kafka一直重试
ack.acknowledge();
}
// 成功才 ack
ack.acknowledge();
}
}
- 1.改进了consumer和producer之间的传输数据的规范性,创建了一个event外壳,里面装入实际传输的对象
- 2.traceId的作用是把一次用户请求在多个服务、多个组件里的日志串起来。一个请求生成一个traceId。
- 3.那traceId用雪花还是什么?traceId一般用uuid,因为traceId对于顺序没有实际的要求,雪花id需要有一定的顺序
- 4.什么是objectMapper?objectMapper 是 Jackson 库里的一个核心类,作用是在Java对象和JSON字符串之间进行转换。
外壳
/**
* kafka事件通用外壳
* @param <T>
*/
@Data
public class EventMessage<T> {
/**
* 消息唯一ID
*/
private String eventId;
/**
* 事件类型
*/
private String eventType;
/**
* 事件版本
*/
private String eventVersion;
/**
* 事件发生时间
*/
private LocalDateTime occurredAt;
/**
* 链路追踪ID
*/
private String traceId;
/**
* 消息来源
*/
private String source;
/**
* 具体业务数据
*/
private T data;
}
8.8 producer
第一版
@Slf4j
@Service
@RequiredArgsConstructor
public class SeckillServiceImpl extends ServiceImpl<SeckillOrderMapper, SeckillOrder> implements SeckillService {
private final StringRedisTemplate redisTemplate;
private final ObjectMapper objectMapper;
// KEY -> VALUE
private final KafkaTemplate<String, String> kafkaTemplate;
private static final String TOPIC = "seckill-order-created";
private final DefaultRedisScript<Long> buildScript;
/**
* 优惠券秒杀
* @param seckillReqVo
* @return
*
*/
@Override
public Long doSeckill(SeckillReqVo seckillReqVo) {
// 1.取出字段拼接redis的key 并 给lua传入参数
long userId = seckillReqVo.getUserId();
long voucherId = seckillReqVo.getVoucherId();
long orderId = IdUtil.getSnowflakeNextId();
String stockKey = "seckill:stock:" + voucherId;
String userSetKey = "seckill:order:users:" + voucherId;
// 2.调用lua脚本
Long result = redisTemplate.execute(buildScript,
Arrays.asList(stockKey, userSetKey),
String.valueOf(userId),
String.valueOf(orderId));
// 3.判断结果
int code = result.intValue();
if(code == 1) {
throw new BusinessException(10001, "库存不足");
} else if(code == 2) {
throw new BusinessException(10002, "重复下单");
} else if(code != 0) {
throw new BusinessException(code, "lua返回异常");
}
// 4.结果正常发送消息
SeckillOrderCreatedEvent event = new SeckillOrderCreatedEvent();
event.setUserId(userId);
event.setVoucherId(voucherId);
event.setOrderId(orderId);
// 5.将消息装入通用事件外壳
EventMessage<SeckillOrderCreatedEvent> eventMessage = new EventMessage();
eventMessage.setData(event);
eventMessage.setEventId(String.valueOf(orderId));
eventMessage.setEventType("SECKILL_ORDER_CREATED");
eventMessage.setSource("SECKILL_ORDER");
eventMessage.setOccurredAt(LocalDateTime.now());
eventMessage.setEventVersion("1.0");
eventMessage.setTraceId(UUID.randomUUID().toString());
String json = "";
try {
json = objectMapper.writeValueAsString(eventMessage);
} catch (Exception e) {
log.info("json序列化失败, orderId = {}", orderId);
}
// TOPIC KEY VALUE
kafkaTemplate.send(KafkaConstants.SECKILL_TOPIC, String.valueOf(userId), json);
return orderId;
}
}
- 注意序列化那一块,万一序列化失败,还是把消息给发出去了,发的还是空字符串,所以需要改进
第二版
@Slf4j
@Service
@RequiredArgsConstructor
public class SeckillServiceImpl extends ServiceImpl<SeckillOrderMapper, SeckillOrder> implements SeckillService {
private final StringRedisTemplate redisTemplate;
private final ObjectMapper objectMapper;
// KEY -> VALUE
private final KafkaTemplate<String, String> kafkaTemplate;
private final DefaultRedisScript<Long> buildScript;
/**
* 优惠券秒杀
* @param seckillReqVo
* @return
*
*/
@Override
public Long doSeckill(SeckillReqVo seckillReqVo) {
// 1.取出字段拼接redis的key 并 给lua传入参数
long userId = seckillReqVo.getUserId();
long voucherId = seckillReqVo.getVoucherId();
long orderId = IdUtil.getSnowflakeNextId();
String stockKey = "seckill:stock:" + voucherId;
String userSetKey = "seckill:order:users:" + voucherId;
// 2.调用lua脚本
Long result = redisTemplate.execute(buildScript,
Arrays.asList(stockKey, userSetKey),
String.valueOf(userId),
String.valueOf(orderId));
// 3.判断结果
int code = result.intValue();
if(code == 1) {
throw new BusinessException(10001, "库存不足");
} else if(code == 2) {
throw new BusinessException(10002, "重复下单");
} else if(code != 0) {
throw new BusinessException(code, "lua返回异常");
}
// 4.结果正常发送消息
SeckillOrderCreatedEvent event = new SeckillOrderCreatedEvent();
event.setUserId(userId);
event.setVoucherId(voucherId);
event.setOrderId(orderId);
// 5.将消息装入通用事件外壳
EventMessage<SeckillOrderCreatedEvent> eventMessage = new EventMessage();
eventMessage.setData(event);
eventMessage.setEventId(String.valueOf(orderId));
eventMessage.setEventType("SECKILL_ORDER_CREATED");
eventMessage.setSource("SECKILL_ORDER");
eventMessage.setOccurredAt(LocalDateTime.now());
eventMessage.setEventVersion("1.0");
eventMessage.setTraceId(UUID.randomUUID().toString());
try {
String json = objectMapper.writeValueAsString(eventMessage);
// TOPIC KEY VALUE
kafkaTemplate.send(KafkaConstants.SECKILL_TOPIC, String.valueOf(userId), json);
} catch (Exception e) {
log.info("json序列化失败或kafka消息发送失败, orderId = {}", orderId);
throw new BusinessException(10003, "json序列化失败");
}
return orderId;
}
}
9.优化
- 秒杀请求先在 Redis 中通过 Lua 脚本完成库存原子预扣和一人一单校验,再投递 Kafka 进行异步削峰。订单最终以数据库落库为准。对于消息发送失败、消费失败、重复消费等场景,通过发送结果确认、消费者幂等校验、失败重试和补偿任务保证最终一致性。
- 为了防止 Kafka 重复投递导致重复创建订单,我在数据库层面对 user_id 和 voucher_id 建立唯一索引,消费者在落库时即使重复执行,也只能成功一次,从而保证消费幂等。
10.如何保证redis和db数据一致
- 普通缓存一致性
我们项目中采用的是 Cache Aside 模式,查询时先查 Redis,未命中再查 MySQL 并回填缓存;更新时采用先更新数据库、再删除缓存的策略,保证缓存和数据库达到最终一致性。
- 秒杀高并发场景
对于秒杀库存这类高并发场景,我们不会直接依赖数据库做实时扣减,而是通过 Redis + Lua 脚本先完成库存校验和扣减,再通过 Kafka 异步落库到 MySQL。这样既保证了高并发性能,又通过消息队列重试和幂等控制实现最终一致性。
- 数据库分库分表优化
当前项目量级下还未引入分库分表,因为过早拆分会提高复杂度;如果后续订单量持续增长,可以按 user_id 或 order_id 做水平分表,并结合雪花算法保证分布式唯一 ID。水平分表其实就是一句话:把一张表的数据按行拆成多张结构相同的表。
- 怎么保证幂等性 什么是幂等性
同一个请求执行一次和执行多次,结果是一样的。这就叫幂等性。如果不保证幂等性的话,会导致kafka重复消费,会一人多单,我使用的是唯一索引来自动保证幂等,在扣库存之前用sismember来判断集合里面是否有已经下单的userId。
- 为什么Redis保证了一人一单,数据库还要再做幂等?
Redis 只能保证请求阶段的幂等,但无法保证消息消费阶段不重复。消息可能被重复消费,redis只保证了前半部分的幂等。
- 如果 Redis 扣库存成功,但 MQ 发送失败怎么办?
这是 Redis 和 MQ 之间的一致性问题。因为 Redis 操作和 Kafka 发送不在同一个事务里,所以无法天然保证强一致。工程上通常会通过发送失败重试、补偿机制、对账任务,或者更进一步采用 Redis Stream 这种能把扣库存和写消息放进同一个 Lua 脚本里的方案来保证最终一致性。对于 Kafka 方案,一般重点是通过重试和补偿保证消息尽量不丢,再通过消费者幂等保证最终落库正确。生产者确认消息最终没有发出去后,才补偿 Redis,把扣除的库存和把该用户从userIdSet去除。用send.get等待一下kafka的回执
- 如果消费者宕机,redis库存已经扣完,但订单还没落库怎么办?
如果消费者宕机,而 Redis 库存已经扣减、Kafka 消息也已经发送成功,这时候不能立即补偿 Redis,因为消息已经进入 Kafka,只是暂时还没有被消费。正确做法是依赖 Kafka 的消息持久化和重试机制,等消费者恢复后继续消费并完成订单落库。同时消费端需要通过唯一索引等方式保证幂等,避免重复创建订单。只有当消息经过多次重试后仍然失败,并被判定为最终无法处理时,才会考虑对 Redis 做补偿。消息已经写进 topic 里,只要还没过保留期,就还在。
- 那 Kafka 消息积压太多怎么办?
Kafka 消息积压太多,说明生产速度大于消费速度。解决思路就是:先定位原因,再从消费能力、分区设计、消息处理耗时、削峰限流这几个方向优化。Kafka 积压太多,就从“扩消费者、加分区、优消费、做监控”四个方向处理。
在优惠券秒杀场景中,如果 Kafka 出现消息堆积,我会先看消费者组 Lag 和消费耗时, Lag就是积压量 offset生产到的数量减去消费到的数量。判断是消费实例不足、分区数不足,还是 订单落库链路变慢。如果是瞬时流量洪峰,我会先通过扩容消费者、增加分区、批量消费等方式快速削峰;如果瓶颈在下游 MySQL 落库,我会进一步优化批量写入、缩短事务、隔离非核心逻辑。对于消费失败的消息,会通过重试和死信队列隔离,避免异常消息阻塞整个分区消费。
长期上会结合监控告警、容量评估和链路异步拆分来防止再次堆积。
- 如何查看Lag
kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group your-group

- kafka如何保证数据不丢和不重复
数据不丢:设置acks=all + retries + enbaenable.idempotence = true,acks=all表示等Leader和ISR副本都写入成功才算成功,retries发送失败自动重试,enbaenable.idempotence防止重试导致重复写入。kafka本身有Broker:多副本机制,
发送消息流程:
1. Producer → Leader(写入)
2. Leader 写入本地日志
3. Follower 从 Leader 拉取数据
4. Follower 写入本地
5. Leader 收到确认
6. 返回成功(acks=all)
重复:设置数据库唯一索引unique,保证幂等性。消费者通过手动提交 offset,保证只有业务处理成功才确认消费。
- 在我的秒杀系统中,Kafka通过多副本 + ISR机制实现自动故障恢复,
在正常情况下Broker宕机可以自动完成Leader切换,无需人工干预。
但如果ISR不足且未开启unclean leader election,则分区会不可用,因为备份的Follower有可能是旧数据,需要人工恢复副本。
因此我在设计中优先保证数据一致性,而不是强可用。