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. ProducerLeader(写入)
2. Leader 写入本地日志
3. FollowerLeader 拉取数据
4. Follower 写入本地
5. Leader 收到确认
6. 返回成功(acks=all)

重复:设置数据库唯一索引unique,保证幂等性。消费者通过手动提交 offset,保证只有业务处理成功才确认消费。

  • 在我的秒杀系统中,Kafka通过多副本 + ISR机制实现自动故障恢复,
    在正常情况下Broker宕机可以自动完成Leader切换,无需人工干预。
    但如果ISR不足且未开启unclean leader election,则分区会不可用,因为备份的Follower有可能是旧数据,需要人工恢复副本。
    因此我在设计中优先保证数据一致性,而不是强可用。
© 版权声明

相关文章