预见猿份
主题
首页面试题在线工具关于我们老苗一对一私教学员评价
实战项目
项目前置基础创新WMS项目Java微服务框架与实战云岚到家项目闪聚支付项目学成在线项目青橙电商项目JVM原理与实战调优分布式事务专题Java高频面试题MySQL从入门到精通Java数据结构与算法老苗一对一私教学员评价blog
blog
  • Java微服务框架与实战

    • 内容介绍
    • 项目安装配置
    • day01-微服务服务注册与发现
    • day01-作业参考
    • day02-服务保护、分布式事务
    • day03-微服务网关与配置中心
    • day04-MQ基础
    • day05-MQ高级
    • day05-RabbitMQ集群部署
    • day06-Elasticsearch入门
    • day07-Elasticsearch高级
    • day08-微服务阶段面试篇上
    • day08-配置同步环境
    • day09-微服务阶段面试篇下
    • day09-redis集群
    • day10-微服务阶段面试题







----- 到底线了 -----

×

欢迎来到预见猿份,本站项目均为站长原创,学习中有问题可直接提交给站长老苗解决(微信:mrt_0607)。

苗润土老师,20余年一线项目经验,2014年加入黑马,星辰wms、云岚到家、学成在线项目作者,历任高级讲师、教学主管及课程研究员。 b站老苗

day05 MQ高级 ​

学习目标 ​

  1. 能够测试生产者重试机制
  2. 能够测试生产者确认机制
  3. 能够说出生产者确认机制的两种方式
  4. 能够说出发送失败处理机制
  5. 能够说出消息持久化机制
  6. 能够测试消费者确认机制
  7. 能够说出消费失败重试机制
  8. 能够说出MQ消息幂等性方案
  9. 能够说出延迟消息方案
  10. 能够实现自动取消超时未支付订单功能
  11. 能够说出保证消息的可靠性的完整方案

B站视频链接

1 消息可靠性 ​

1.1. 思路分析 ​

在昨天的练习作业中,我们改造了余额支付功能,在支付成功后利用RabbitMQ通知交易服务,更新业务订单状态为已支付。

但是大家思考一下,如果这里MQ通知失败,支付服务中支付流水显示支付成功,而交易服务中的订单状态却显示未支付,数据出现了不一致。

画板0

首先,我们一起分析一下消息丢失的可能性有哪些。

消息从发送者发送消息,到消费者处理消息,需要经过的流程是这样的:

画板1

消息从生产者到消费者的每一步都可能导致消息丢失:

  • 发送消息时丢失:

    • 生产者发送消息时连接MQ失败
    • 生产者发送消息到达MQ后未找到Exchange
    • 生产者发送消息到达MQ的Exchange后,未找到合适的Queue
    • 消息到达MQ后,处理消息的进程发生异常
  • MQ导致消息丢失:

    • 消息到达MQ,保存到队列后,尚未消费就突然宕机
  • 消费者处理消息时:

    • 消息接收后尚未处理突然宕机
    • 消息接收后处理过程中抛出异常

综上,我们要解决消息丢失问题,保证MQ的可靠性,就必须从3个方面入手:

  1. 保证生产消息的可靠性
  2. 确保MQ不会将消息弄丢
  3. 保证消费消息的可靠性

注意:使用MQ并不是所有场景对消息的可靠性要求都很高,比如上图中,支付成功短信通知的流程对消息可靠性要求就不高,通常都可以保证消息正常到达消费者,即使个别没有成功通知用户也不影响主体业务流程,所以在设计技术方案时一定要根据业务需求具体分析。

1.2 生产消息可靠性 ​

1.2.1. 生产者重试机制 ​

1.2.1.1 配置 ​

首先第一种情况,就是生产者发送消息时,出现了网络故障,导致与MQ的连接中断。

为了解决这个问题,SpringAMQP提供的消息发送时的重试机制。即:当RabbitTemplate与MQ连接超时后,多次重试。

从课程资料中找到mq-demo-v2.zip,并解压,使用IDEA打开mq-demo-v2工程。

修改publisher模块的application.yaml文件,添加下面的内容:

YAML
spring:
  rabbitmq:
    connection-timeout: 1s # 设置MQ的连接超时时间
    template:
      retry:
        enabled: true # 开启超时重试机制
        initial-interval: 1000ms # 失败后的初始等待时间
        multiplier: 1 # 失败后下次的等待时长倍数,下次等待时长 = 上次等待时长 * multiplier
        max-attempts: 3 # 总共尝试次数
1
2
3
4
5
6
7
8
9

配置参数解释

  • initial-interval: 失败后的初始等待时间
  • multiplier: 倍增器,每次重试的等待时间是前一次几倍。

失败后下次等待时长 =上次等待时长 * multiplier

  • max-attempts: 最大重试次数(包括第一次尝试)

举例:

由于multiplier设置为1,这意味着每次重试之间的间隔是固定的,不会增加。

假设在t=0时刻首次尝试发送消息,如果发送失败,则会按照以下时间点进行重试:

第一次尝试(也是首次发送):t=0(假设即时失败)

第一次重试:等待1秒后重试,t=1秒(首次失败后等待1秒)

第二次重试:等待1秒*1=1秒 后重试,t=2秒(从第一次重试再等待1秒)

如果设置如下:

JavaScript
initial-interval:1000ms
multiplier:2
max-attempts: 5
1
2
3

由于multiplier设置为2,这意味着每次重试之间的间隔会翻倍。

假设在t=0时刻首次尝试发送消息,如果发送失败,则会按照以下时间点进行重试:

第一次尝试(也是首次发送):t=0(假设即时失败)

第一次重试:等待1秒后重试,t=1秒(首次失败后等待1秒)

第二次重试:等待1*2=2秒 后重试,t=3秒

第三次重试:等待2*2=4秒 后重试,t=7秒

第四次重试:等待4*2=8秒 后重试,t=15秒

1.2.1.2 测试 ​

我们利用命令停掉RabbitMQ服务:

Shell
docker stop mq
1

然后测试发送一条消息,会发现会每隔1秒重试1次,总共重试了3次。消息发送的超时重试机制配置成功了!

注意:当网络不稳定的时候,利用重试机制可以有效提高消息发送的成功率。不过SpringAMQP提供的重试机制是阻塞式的重试,也就是说多次重试等待的过程中,当前线程是被阻塞的。

如果对于业务性能有要求,建议禁用重试机制。如果一定要使用,请合理配置等待时长和重试次数,当然也可以考虑使用异步线程来执行发送消息的代码。

1.2.2. 生产者确认机制 ​

1.2.2.1 两种机制介绍 ​

一般情况下,只要生产者与MQ之间的网路连接顺畅,基本不会出现发送消息丢失的情况,因此大多数情况下我们无需考虑这种问题。

不过,在少数情况下,也会出现消息发送到MQ之后丢失的现象,比如:

  • MQ内部处理消息的进程发生了异常
  • 生产者发送消息到达MQ后未找到Exchange
  • 生产者发送消息到达MQ的Exchange后,未找到合适的Queue,因此无法路由

针对上述情况,RabbitMQ提供了生产者消息确认机制,包括Publisher Confirm和Publisher Return两种。在开启确认机制的情况下,当生产者发送消息给MQ后,MQ会根据消息处理的情况返回不同的回执。

具体如图所示:

配图00

生产者确认机制:

1.Publisher Return

消息投递成功但路由失败会调用Publisher Return回调方法返回异常信息。

2.Publisher Confirm

消息投递成功返回ack,投递失败返回nack。

注意:消息投递成功但可能路由失败了,此时会通过Publisher Confirm返回ack,通过Publisher Return回调方法返回异常信息。

默认两种机制都是关闭状态,需要通过配置文件来开启。

1.2.2.2 开启生产者确认 ​

在publisher模块的application.yaml中添加配置:

YAML
spring:
  rabbitmq:
    publisher-confirm-type: correlated # 开启publisher confirm机制,并设置confirm类型
    publisher-returns: true # 开启publisher return机制
1
2
3
4

这里publisher-confirm-type有三种模式可选:

  • none:关闭confirm机制
  • simple:同步阻塞等待MQ的回执(回调方法)
  • correlated:MQ异步回调返回回执

一般我们推荐使用correlated,回调机制。

1.2.2.3 实现方法 ​

ReturnCallback ​

每个RabbitTemplate只能配置一个ReturnCallback,因此我们可以在配置类中统一设置。我们在publisher模块定义一个配置类:

配图01

内容如下:

Java
package com.itheima.publisher.config;

import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.ReturnedMessage;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.context.annotation.Configuration;

import javax.annotation.PostConstruct;

@Slf4j
@AllArgsConstructor
@Configuration
public class MqConfig {
    private final RabbitTemplate rabbitTemplate;

    @PostConstruct
    public void init(){
        rabbitTemplate.setReturnsCallback(new RabbitTemplate.ReturnsCallback() {
            @Override
            public void returnedMessage(ReturnedMessage returned) {
                log.error("触发return callback,");
                log.debug("exchange: {}", returned.getExchange());
                log.debug("routingKey: {}", returned.getRoutingKey());
                log.debug("message: {}", returned.getMessage());
                log.debug("replyCode: {}", returned.getReplyCode());
                log.debug("replyText: {}", returned.getReplyText());
            }
        });
    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
ConfirmCallback ​

由于每个消息发送时的处理逻辑不一定相同,因此ConfirmCallback需要在每次发消息时定义。具体来说,是在调用RabbitTemplate中的convertAndSend方法时,多传递一个参数:

配图02

这里的CorrelationData中包含两个核心的东西:

  • id:消息的唯一标示,MQ对不同的消息的回执以此做判断,避免混淆
  • SettableListenableFuture:回执结果的Future对象

将来MQ的回执就会通过这个Future来返回,我们可以提前给CorrelationData中的Future添加回调函数来处理消息回执:

配图03

我们测试下边的方法,向系统自带的交换机发送消息,并且添加ConfirmCallback:

注意:此代码不用编写直接测试即可,稍后我们会用工具类替代。

Java
@Test
void testPublisherConfirm() {
    // 1.创建CorrelationData
    CorrelationData cd = new CorrelationData();
    // 2.给Future添加ConfirmCallback
    cd.getFuture().addCallback(new ListenableFutureCallback() {
        @Override
        public void onFailure(Throwable ex) {
            // 2.1.Future发生异常时的处理逻辑,基本不会触发
            log.error("send message fail", ex);
        }
        @Override
        public void onSuccess(CorrelationData.Confirm result) {
            // 2.2.Future接收到回执的处理逻辑,参数中的result就是回执内容
            if(result.isAck()){ // result.isAck(),boolean类型,true代表ack回执,false 代表 nack回执
                log.debug("发送消息成功,收到 ack!");
            }else{ // result.getReason(),String类型,返回nack时的异常描述
                log.error("发送消息失败,收到 nack, reason : {}", result.getReason());
            }
        }
    });
    // 3.发送消息,故意指定一个错误的rontingKey
    rabbitTemplate.convertAndSend("hmall.direct", "q", "hello", cd);
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24

1.2.2.4 测试 ​

测试报错:Initialization of bean failed; nested exception is java.lang.IllegalStateException: Only one ReturnCallback is supported by each RabbitTemplate

原因:每个RabbitTemplate只支持一个ReturnCallback。

解决:

屏蔽common-rabbitmq依赖,在common-rabbitmq中对RabbitTemplate设置了ReturnCallback

配图04

测试步骤:

可以看到,由于传递的RoutingKey是错误的,路由失败后,触发了return callback,同时也收到了ack。

配图05

当我们修改为正确的RoutingKey以后,就不会触发return callback了,只收到ack。

当我们把交换机名称修改错误则只会收到nack。

1.2.3. 发送失败处理机制 ​

1.2.3.1 失败处理机制 ​

在ConfirmCallback中收到nack表示消息投递失败,ReturnCallback异常表示路由失败,消息投递失败怎么处理?

可以将消息记录到失败消息表,由定时任务进行发布,每隔10秒钟(可设置)执行获取失败消息重新发送,发送一次则在失败次数字段加一,达到3次停止自动发送由人工处理。

在commonn-rabbitmq模块中实现了发送消息的工具方法,此方法实现了发送失败处理机制。

ReturnCallback回调逻辑在com.itheima.common.rabbitmq.config.RabbitMqConfiguration中,核心代码如下:

Java
@Configuration
@ConditionalOnProperty(prefix = "rabbit-mq", name = "enable", havingValue = "true")
@Import({RabbitClient.class, FailMsgDaoImpl.class})
@Slf4j
public class RabbitMqConfiguration implements ApplicationContextAware {

    /**
     * 并发数量
     */
    public static final int DEFAULT_CONCURRENT = 10;

    @Autowired(required = false)
    private FailMsgDao failMsgDao;


    @Override
    public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
        // 获取RabbitTemplate
        RabbitTemplate rabbitTemplate = applicationContext.getBean(RabbitTemplate.class);
        //定义returnCallback回调方法
        rabbitTemplate.setReturnsCallback(
                new RabbitTemplate.ReturnsCallback() {
                    @Override
                    public void returnedMessage(ReturnedMessage returnedMessage) {
                        byte[] body = returnedMessage.getMessage().getBody();
                        //消息id
                        String messageId = returnedMessage.getMessage().getMessageProperties().getMessageId();
                        String content = new String(body, Charset.defaultCharset());
                        log.info("消息发送失败,应答码{},原因{},交换机{},路由键{},消息id{},消息内容{}",
                                returnedMessage.getReplyCode(),
                                returnedMessage.getReplyText(),
                                returnedMessage.getExchange(),
                                returnedMessage.getRoutingKey(),
                                messageId,
                                content);
                        if (failMsgDao != null) {
                            failMsgDao.save(messageId, returnedMessage.getExchange(), returnedMessage.getRoutingKey(), content, 0, DateUtils.getCurrentTime()+10, "returnCallback");
                        }
                    }
                }
        );
    }

}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44

ApplicationContextAware 的作用:如果 Bean 实现了 ApplicationContextAware 接口,Spring 容器会调用 setApplicationContext 方法,将 ApplicationContext 传递给该 Bean。

Bean 创建:Spring 容器首先创建 Bean 实例。

属性注入:Spring 容器对 Bean 的属性进行依赖注入。

Aware 接口回调:如果 Bean 实现了 ApplicationContextAware 接口,Spring 容器会调用 setApplicationContext 方法,将 ApplicationContext 传递给该 Bean。

初始化方法:如果 Bean 配置了初始化方法(例如通过 @PostConstruct 注解或 init-method 属性),Spring 容器会调用这些初始化方法。

ConfirmCallback的逻辑在RabbitClient 工具类的sendMsg方法中,通过sendMsg方法发送消息,在发送消息时指定

correlationData对象,在correlationData对象中指定了ConfirmCallback回调方法的逻辑,当返回nack会将消息写入失败表,如果消息重发成功会将该记录从失败表删除。

核心代码如下:

JavaScript
@Slf4j
@Service
public class RabbitClient {
    ....

    @Resource
    private RabbitTemplate rabbitTemplate;
    @Autowired(required = false)
    private FailMsgDao failMsgDao;
    
    /**
 * 发送消息 重试3次
 *
 * @param exchange   交换机
 * @param routingKey 路由key
 * @param msg        消息对象,会将对象序列化成json字符串发出
 * @param delay      延迟时间 秒
 * @param msgId      消息id
 * @param isFailMsg  是否是失败消息
 * @return 是否发送成功
 */
@Retryable(value = MqException.class, maxAttempts = 3, backoff = @Backoff(value = 3000, multiplier = 1.5), recover = "saveFailMag")
public void sendMsg(String exchange, String routingKey, Object msg, Integer delay, String msgId, boolean isFailMsg) {
    // 1.发送消息前准备
    // 1.1获取消息内容,如果非字符串将其序列化
    String jsonMsg = JsonUtils.toJsonStr(msg);
    // 1.2.全局唯一消息id,如果调用者设置了消息id,使用调用者消息id,如果为配置,默认雪花算法生成消息id
    if(StrUtil.isBlank(msgId)){
        msgId = IdUtil.getSnowflakeNextIdStr();
    }
    // 1.3.设置默认延迟时间,默认立即发送
    delay = NumberUtils.null2Default(delay, -1);
    log.debug("消息发送!exchange = {}, routingKey = {}, msg = {}, msgId = {}", exchange, routingKey, jsonMsg, msgId);


    // 1.4.构建回调
    RabbitMqListenableFutureCallback futureCallback = RabbitMqListenableFutureCallback.builder()
            .exchange(exchange)
            .routingKey(routingKey)
            .msg(jsonMsg)
            .msgId(msgId)
            .delay(delay)
            .isFailMsg(isFailMsg)
            .failMsgDao(failMsgDao)
            .build();
    // 1.5.CorrelationData设置
    CorrelationData correlationData = new CorrelationData(msgId.toString());
    correlationData.getFuture().addCallback(futureCallback);

    // 1.6.构造消息对象
    Message message = MessageBuilder.withBody(StrUtil.bytes(jsonMsg, CharsetUtil.CHARSET_UTF_8))
            //持久化
            .setDeliveryMode(MessageDeliveryMode.PERSISTENT)
            //消息id
            .setMessageId(msgId.toString())
            .build();

    try {
        // 2.发送消息
        this.rabbitTemplate.convertAndSend(exchange, routingKey, message, new DelayMessagePostProcessor(delay), correlationData);
    } catch (Exception e) {
        log.error("send error:" + e);
        // 3.构建异常回调,并抛出异常
        MqException mqException = new MqException();
        mqException.setMsg(ExceptionUtil.getMessage(e));
        mqException.setMqId(msgId);
        throw mqException;
    }
}
    
 public class RabbitMqListenableFutureCallback implements ListenableFutureCallback {

    //记录失败消息service
    private FailMsgDao failMsgDao;

    private String exchange;
    private String routingKey;
    private String msg;
    private String msgId;
    private Integer delay;

    //是否是失败消息
    private boolean isFailMsg=false;

    @Override
    public void onFailure(Throwable ex) {
        if(failMsgDao == null) {
            return;
        }
        failMsgDao.save(msgId, exchange, routingKey, msg, delay, DateUtils.getCurrentTime() + 10, ExceptionUtil.getMessage(ex));
    }

    @Override
    public void onSuccess(CorrelationData.Confirm result) {
        if(failMsgDao == null){
            return;
        }
        if(!result.isAck()){
            // 执行失败保存失败信息,如果已经存在保存信息,如果不在信息信息
            failMsgDao.save(msgId, exchange, routingKey, msg, delay,DateUtils.getCurrentTime() + 10, "MQ回复nack");
        }else if(isFailMsg && msgId != null){
            // 如果发送的是失败消息,当收到ack需要从fail_msg删除该消息
            failMsgDao.removeById(msgId);
        }
    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106

当发送消息失败会入库到失败消息表。

我们可以启动定时任务去扫描失败消息表的记录,重新发送,当达到最大失败次数后由人工处理。

1.2.3.4 测试 ​

下边测试发送失败入库功能:

  1. 首先在publisher模块添加common-rabbitmq依赖
  2. 屏蔽MqConfig类中设置ReturnCallback的代码

由于RabbitTemplate只能设置一次ReturnCallback,而在common-rabbitmq中设置了ReturnCallback,所以屏蔽MqConfig类中设置ReturnCallback的代码,如下图:

配图06

  1. 在hmall中创建失败消息表

查看hmall数据库,如果已存在fail_msg表需要删除原表再通过下边的语句新建表。

Java
create table fail_msg
(
    id                     varchar(255)  not null comment '消息id'
        primary key,
    exchange               varchar(255)  not null comment '交换机',
    routing_key            varchar(255)  not null comment '路由key',
    msg                    text          not null comment '消息',
    reason                 varchar(255)  not null comment '原因',
    delay_msg_execute_time int           null comment '延迟消息执行时间',
    next_fetch_time        int           null comment '下次拉取时间',
    create_time            datetime      null comment '创建时间',
    update_time            datetime      null comment '更新时间',
    fail_count             int default 0 null comment '失败次数'
)
    comment '失败消息记录表' charset = utf8mb4;
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
  1. 在publisher模块进行配置

在启动类中添加@MapperScan("com.itheima.common.rabbitmq.dao.mapper")

在application.yaml添加

Java
rabbit-mq:
  enable: true
  persistence:
    enable: true
1
2
3
4
  1. 使用RabbitClient工具类发送消息
Java
@Resource
private RabbitClient rabbitClient;
@Test
void testPublisherReturn() {
    rabbitClient.sendMsg("hmall.directa", "q", "hello");
}
1
2
3
4
5
6

分别测试ReturnCallback和ConfirmCallback不同情况下失败消息入库。

测试ReturnCallback时注意:单元测试方法运行完就完毕了数据库连接池,而ReturnCallback是回调方法,是在单元测试方法执行完再执行,在ReturnCallback中操作数据库时报没有可用的数据库连接的错误,需要在单元测试方法最后添加休眠代码,保证ReturnCallback执行完成再结束整个单元测试方法。

1.2.3 小结 ​

如何保证生产消息可靠性?

首先在发送消息时可以开启重试机制,避免因为短暂的网络问题导致发送消息失败。

RabbitMQ还提供生产者确认机制保证发送消息到MQ的可靠性。

生产者确认机制包括两种:

1.Publisher Return

消息投递成功但路由失败会调用Publisher Return回调方法返回异常信息。

2.Publisher Confirm

消息投递成功返回ack,投递失败返回nack。

注意:消息投递成功但可能路由失败了,此时会通过Publisher Confirm返回ack,通过Publisher Return回调方法返回异常信息。

我们在发送消息时给每个消息指定一个唯一ID,设置回调方法,如果Publisher Return失败或Publisher Confirm返回nack,我们在回调方法中解析失败消息,并记录到失败表由定时任务去异步重新发送。

1.3.消息持久化 ​

为了提升性能,默认情况下MQ的数据都是在内存存储的临时数据,重启后就会消失。为了保证数据的可靠性,必须配置持久化,包括:

  • 交换机持久化
  • 队列持久化
  • 消息持久化

我们以控制台界面为例来说明。

1.3.1 交换机持久化 ​

交换机持久化是指将交换机的定义信息(元数据)持久化到RabbitMQ的数据库(mnesia)中,RabbitMQ重启后交换机定义仍然存在。

若交换机不设置持久化,在rabbitmq服务重启之后,相关的交换机元数据会丢失,但消息不会丢失,只是不能将消息发送到这个交换机中,所以通常要设置交换机持久化。

在控制台中声明交换机是默认设置持久化的。

设置方法:

在控制台的Exchanges页面,添加交换机时可以配置交换机的Durability参数:

配图07

设置为Durable就是持久化模式,Transient就是临时模式。

设置持久化后在交换机列表会有一个"D"标识

配图08

1.3.2 队列持久化 ​

队列持久化也是将队列的定义信息(元数据)持久化到RabbitMQ的数据库中。

如果队列不设置持久化,在RabbitMQ重启后队列的元数据丢失。

设置队列持久化可以保证队列本身的元数据不会因异常情况而丢失,队列中存储的是消息的在队列中的位置、消息的ID、存储位置等,消息会存储在独立的rdq数据文件中,队列持久化不能保证消息数据不会丢失。

在控制台的Queues页面,添加队列时,同样可以配置队列的Durability参数:

配图09

除了持久化以外,你可以看到队列还有很多其它参数,有一些我们会在后期学习。

设置持久化后在队列列表会有一个"D"标识

配图10

1.3.3 消息持久化 ​

要确保消息不会丢失需要将消息设置为持久化。

在控制台发送消息的时候,可以添加很多参数,而消息的持久化是要配置一个properties:

配图11

查看持久消息的内容:

delivery_mode=2 表示持久化

配图12

1.4 消费消息可靠性 ​

当RabbitMQ向消费者投递消息以后,需要知道消费者的处理状态如何。因为消息投递给消费者并不代表就一定被正确消费了,可能出现的故障有很多,比如:

  • 消息投递的过程中出现了网络故障
  • 消费者接收到消息后突然宕机
  • 消费者接收到消息后,因处理不当导致异常
  • ...

一旦发生上述情况,消息也会丢失。因此,RabbitMQ必须知道消费者的处理状态,一旦消息处理失败能重新投递消息。

但问题来了:RabbitMQ如何得知消费者的处理状态呢?

本章我们就一起研究一下消费者处理消息时的可靠性解决方案。

1.4.1. 消费者确认机制 ​

1.4.1.1 介绍 ​

为了确认消费者是否成功处理消息,RabbitMQ提供了消费者确认机制(Consumer Acknowledgement)。即:当消费者处理消息结束后,应该向RabbitMQ发送一个回执,告知RabbitMQ消息处理状态。回执有三种可选值:

  • ack:成功处理消息,RabbitMQ从队列中删除该消息
  • nack:消息处理失败,RabbitMQ需要再次投递消息
  • reject:消息处理失败并拒绝该消息,RabbitMQ从队列中删除该消息

一般reject方式用的较少,除非是消息格式有问题,那就是开发问题了。因此大多数情况下我们需要将消息处理的代码通过try catch机制捕获,消息处理成功时返回ack,处理失败时返回nack.

SpringAMQP帮我们实现了消息确认,并可以通过配置文件设置消息确认的处理方式,有三种模式:

  • none:不处理。即消息投递给消费者后消息会立刻从MQ删除。非常不安全,不建议使用

  • manual:手动模式。需要自己在业务代码中调用api,发送ack或reject,存在业务入侵,但更灵活

  • auto:自动模式。当业务正常执行时则自动返回ack. 当业务出现异常时,根据异常判断返回不同结果:

    • 如果是业务异常,会自动返回nack;
    • 如果是消息处理或校验异常,自动返回reject,返回的异常包括:MessageConversionException、MethodArgumentTypeMismatchException等

1.4.1.2 auto模式测试 ​

通过下面的配置可以修改消息确认的处理方式为auto:

YAML
spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: auto # 自动ack
1
2
3
4
5

修改consumer服务的SpringRabbitListener类中的方法,模拟一个消息处理的异常:

Java
@RabbitListener(queues = "simple.queue")
public void listenSimpleQueueMessage2(String msg) throws InterruptedException {
    log.info("spring 消费者接收到消息:【" + msg + "】");
    if (true) {
        throw new RuntimeException("故意的");
    }
    log.info("消息处理完成");
}
1
2
3
4
5
6
7
8

在此方法内第一句打断点,我们向队列“simple.queue”发一条消息,此时进入断点

配图13

可以发现此时有一条消息状态为unacked(未确定状态):

配图14

放行以后,由于抛出的是业务异常,消息处理失败后,会回到RabbitMQ,并重新投递到消费者。

1.4.1.3 manual模式测试(自学) ​

通常在应用程序中会使用auto自动模式,手动模式在生产中不常用可以自行学习。

通过 api指定是否重新入队,具体返回ack或nack通过手动编程实现。

由于手动模式需要通过api编程,需要在监听方法添加Channel、Message类型的参数,如下:

Message:是spring AMQP封装的底层消息对象。

Channel:是消费端与MQ基于通道的操作对象。

JavaScript
    @RabbitListener(queues = "simple.queue")
    public void listenSimpleQueueMessage(String msg, Channel channel, Message message) throws InterruptedException, IOException {
        log.info("spring 消费者接收到消息:【" + msg + "】");
        //返回nack
        //每个参数的意义:1.消息的标记 2.是否确认之前所有未确认的消息 3.是否重新入队
        channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
//        log.info("消息处理完成");
//        //返回ack,每个参数的意义:1.消息的标记 2.是否确认之前所有消息
//        channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
    }
1
2
3
4
5
6
7
8
9
10

1.4.1.4 小结 ​

消息者确认机制怎么实现?

消费者处理消息结束后,应该向RabbitMQ发送一个回执,告知RabbitMQ消息处理状态:

  • ack:成功处理消息,RabbitMQ从队列中删除该消息
  • nack:消息处理失败,RabbitMQ需要再次投递消息
  • reject:消息处理失败并拒绝该消息,RabbitMQ从队列中删除该消息

具体实现方法:

SpringAMQP提供消息确认配置,有三种模式:

  • none:不处理。即消息投递给消费者后消息会立刻从MQ删除。非常不安全,不建议使用。
  • manual:手动模式。需要自己在业务代码中调用api,发送ack或reject。
  • auto:自动模式。当业务正常执行时则自动返回ack. 当业务出现异常时会自动返回nack.

我们通常使用auto自动模式,业务处理完成没有异常则自动返回ack,如果存在业务异常则自动返回nack,如果在消息处理或消息校验异常时自动返回reject。

除了自动确认还要了解手动确认,手动确认是将返回ack或nack在程序中编码实现,举例如下:

返回ack:

JavaScript
    //返回ack,每个参数的意义:1.消息的标记 2.是否确认之前所有消息
    channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
1
2

返回nack:

JavaScript
 //每个参数的意义:1.消息的标记 2.是否确认之前所有消息 3.是否重新入队
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, false);
1
2

1.4.2. 失败重试机制 ​

1.4.2.1 本地重试机制 ​

当消费者出现异常后,消息会不断requeue(重入队)到队列,再重新发送给消费者。如果消费者再次执行依然出错,消息会再次返回到队列,再次投递,直到消息处理成功为止。

极端情况就是消费者一直无法执行成功,那么消息投递就会无限循环,导致mq的消息处理飙升,带来不必要的压力:

配图15

当然,上述极端情况发生的概率还是非常低的,不过不怕一万就怕万一。为了应对上述情况Spring又提供了消费者失败重试机制:在消费者出现异常时利用本地重试,而不是无限制的投递到mq队列。

修改consumer服务的application.yml文件,添加内容:

YAML
spring:
  rabbitmq:
    listener:
      simple:
        retry:
          enabled: true # 开启消费者失败重试
          initial-interval: 1000ms # 初识的失败等待时长为1秒
          multiplier: 1 # 失败的等待时长倍数,下次等待时长 = 上次等待时长 * multiplier
          max-attempts: 3 # 最大重试次数
1
2
3
4
5
6
7
8
9

重启consumer服务,重复之前的测试。可以发现:

  • 消费者在失败后消息没有重新回到MQ无限重新投递,而是在本地重试了3次
  • 本地重试3次以后,查看RabbitMQ控制台,发现消息被删除了,说明最后SpringAMQP返回的是reject

结论:

  • 开启本地重试时,消息处理过程中抛出异常,不会请求到队列,而是在消费者本地重试
  • 重试达到最大次数后,Spring会返回reject,消息会被丢弃

1.4.2.2 失败消息入队 ​

本地测试达到最大重试次数后,消息会被丢弃。这在某些对于消息可靠性要求较高的业务场景下,显然不太合适了。

因此Spring允许我们自定义重试次数耗尽后的消息处理策略,这个策略是由MessageRecovery接口来定义的,它有3个不同实现:

  • RejectAndDontRequeueRecoverer:重试耗尽后,直接reject,丢弃消息。默认就是这种方式
  • ImmediateRequeueMessageRecoverer:重试耗尽后,返回nack,消息重新入队
  • RepublishMessageRecoverer:重试耗尽后,将失败消息投递到指定的交换机

比较优雅的一种处理方案是RepublishMessageRecoverer,失败后将消息投递到一个固定交换机,通过交换机将消息转发到失败消息队列,程序监听失败消息队列,接收到失败消息,将失败消息存入失败消息表,通过定时任务进行处理。

1)在consumer服务中定义处理失败消息的交换机和队列

Java
@Bean
public DirectExchange errorMessageExchange(){
    return new DirectExchange("error.direct");
}
@Bean
public Queue errorQueue(){
    return new Queue("error.queue", true);
}
@Bean
public Binding errorBinding(Queue errorQueue, DirectExchange errorMessageExchange){
    return BindingBuilder.bind(errorQueue).to(errorMessageExchange).with("error");
}
1
2
3
4
5
6
7
8
9
10
11
12

2)定义一个RepublishMessageRecoverer,指定失败消息投递交换机的名称及routingkey

Java
@Bean
public MessageRecoverer republishMessageRecoverer(RabbitTemplate rabbitTemplate){
    return new RepublishMessageRecoverer(rabbitTemplate, "error.direct", "error");
}
1
2
3
4

完整代码如下:

Java
package com.itheima.consumer.config;

import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.rabbit.retry.MessageRecoverer;
import org.springframework.amqp.rabbit.retry.RepublishMessageRecoverer;
import org.springframework.context.annotation.Bean;

@Configuration
@ConditionalOnProperty(name = "spring.rabbitmq.listener.simple.retry.enabled", havingValue = "true")
public class ErrorMessageConfig {
    @Bean
    public DirectExchange errorMessageExchange(){
        return new DirectExchange("error.direct");
    }
    @Bean
    public Queue errorQueue(){
        return new Queue("error.queue", true);
    }
    @Bean
    public Binding errorBinding(Queue errorQueue, DirectExchange errorMessageExchange){
        return BindingBuilder.bind(errorQueue).to(errorMessageExchange).with("error");
    }

    @Bean
    public MessageRecoverer republishMessageRecoverer(RabbitTemplate rabbitTemplate){
        return new RepublishMessageRecoverer(rabbitTemplate, "error.direct", "error");
    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32

1.4.2.3 测试失败消息入队 ​

测试流程:

启动消费端程序:

JavaScript
@RabbitListener(queues = "simple.queue")
public void listenSimpleQueueMessage(String msg) throws InterruptedException {
    log.info("spring 消费者接收到消息:【" + msg + "】");
    if (true) {
        throw new RuntimeException("故意的");
    }
    log.info("消息处理完成");
}
1
2
3
4
5
6
7
8

达到最大重试次数将会投递到失败消息队列。

监听失败消息队列将失败消息写入数据库中,由人工定期处理。

JavaScript
@RabbitListener(queues = "error.queue")
public void listenErrorQueue(String msg) throws InterruptedException {
    System.out.println("接收失败消息:【" + msg + "】" + LocalTime.now());
    //存入数据库失败消息表...
}
1
2
3
4
5

1.4.2.4 小结 ​

消费消息失败重试机制是什么?

提供三种策略:

  • RejectAndDontRequeueRecoverer:重试耗尽后,直接reject,丢弃消息。默认就是这种方式
  • ImmediateRequeueMessageRecoverer:重试耗尽后,返回nack,消息重新入队
  • RepublishMessageRecoverer:重试耗尽后,将失败消息投递到指定的交换机

推荐使用RepublishMessageRecoverer,将失败消息投递到固定的交换机,通过交换机将消息转发到失败消息队列,程序监听失败消息队列,接收到失败消息,将失败消息存入失败消息表,通过定时任务进行处理。

1.4.3. MQ消息幂等性 ​

1.4.3.1 什么是幂等性 ​

何为幂等性?

在程序开发中,是指同一个业务,执行一次或多次对业务状态的影响是一致的。例如:

  • 根据id删除数据
  • 查询数据

但数据的更新往往不是幂等的,如果重复执行可能造成不一样的后果。比如:

  • 取消订单,恢复库存的业务。如果多次恢复就会出现库存重复增加的情况
  • 退款业务。重复退款对商家而言会有经济损失。

所以,我们要尽可能避免业务被重复执行,然而在实际业务场景中,由于意外经常会出现业务被重复执行的情况。例如:

  • 页面卡顿时频繁刷新导致表单重复提交
  • 服务间调用的重试
  • MQ消息的重复投递

因此,我们必须想办法保证消息处理的幂等性。这里给出两种方案:

  • 唯一消息ID
  • 业务状态判断

1.4.4.2 唯一消息ID ​

这个思路非常简单:

  1. 每一条消息都生成一个唯一的id,与消息一起投递给消费者。
  2. 消费者接收到消息后处理自己的业务,业务处理成功后将消息ID保存到数据库或Redis
  3. 如果下次又收到相同消息,去数据库或Redis查询判断是否存在,存在则为重复消息放弃处理。

我们该如何给消息添加唯一ID呢?

其实很简单,SpringAMQP的MessageConverter自带了MessageID的功能,我们只要开启这个功能即可。

以Jackson的消息转换器为例:

Java
@Bean
public MessageConverter messageConverter(){
    // 1.定义消息转换器
    Jackson2JsonMessageConverter jjmc = new Jackson2JsonMessageConverter();
    // 2.配置自动创建消息id,用于识别不同消息,也可以在业务中基于ID判断是否是重复消息
    jjmc.setCreateMessageIds(true);
    return jjmc;
}
1
2
3
4
5
6
7
8

接收消息的方法中添加org.springframework.amqp.core.Message类型的参数,整体的处理逻辑是:

  1. 接收消息解析出消息id
  2. 根据消息id查询数据库已处理消息表如果存在说明已处理该消息,如果查询不到说明未处理该消息
  3. 根据判断结果去处理消息
  4. 处理消息完毕将消息id写入已处理消息表。

代码如下:

JavaScript
@RabbitListener(queues = "simple.queue")
public void listenSimpleQueueMessage(String msg,Channel channel,Message message) throws InterruptedException {
    log.info("spring 消费者接收到消息:【" + msg + "】");
    //消息id
    String messageId = message.getMessageProperties().getMessageId();
    //先从数据库查询是否已处理该消息,否则已处理直接return
    if (true) {
        throw new RuntimeException("故意的");
    }
    //将已处理的消息id存入数据库中...
    log.info("消息处理完成");
}
1
2
3
4
5
6
7
8
9
10
11
12

1.4.4.3 业务判断 ​

业务判断就是基于业务本身的逻辑或状态来判断是否是重复的请求,不同的业务场景判断的思路也不一样。

例如在支付通知案例中,处理消息的业务逻辑是把订单状态从未支付修改为已支付。因此我们就可以在执行更新时判断订单状态是否是未支付,如果不是则证明订单已经被处理过,无需重复处理。

相比较而言,使用唯一消息ID的方案需要操作数据库或Redis保存消息ID,所以更推荐使用业务判断的方案。

以支付修改订单的业务为例,我们需要修改OrderServiceImpl中的markOrderPaySuccess方法:

Java
    @Override
    public void markOrderPaySuccess(Long orderId) {
        // 1.查询订单
        Order old = getById(orderId);
        // 2.判断订单状态
        if (old == null || old.getStatus() != 1) {
            // 订单不存在或者订单状态不是1,放弃处理
            return;
        }
        // 3.尝试更新订单
        Order order = new Order();
        order.setId(orderId);
        order.setStatus(2);
        order.setPayTime(LocalDateTime.now());
        updateById(order);
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16

直接update语句中添加条件也可以,如下:

Java
@Override
public void markOrderPaySuccess(Long orderId) {
    // UPDATE `order` SET status = ? , pay_time = ? WHERE id = ? AND status = 1
    lambdaUpdate()
            .set(Order::getStatus, 2)
            .set(Order::getPayTime, LocalDateTime.now())
            .eq(Order::getId, orderId)
            .eq(Order::getStatus, 1)
            .update();
}
1
2
3
4
5
6
7
8
9
10

注意看,上述代码等同于这样的SQL语句:

SQL
UPDATE `order` SET status = ? , pay_time = ? WHERE id = ? AND status = 1
1

我们在where条件中除了判断id以外,还加上了status必须为1的条件。如果条件不符(说明订单已支付),则SQL匹配不到数据,根本不会执行。

1.4.4.4 面试题 ​

如何保证MQ幂等性?或 如何防止消息重复消费?

1.4.4. 业务补偿 ​

虽然我们利用各种机制尽可能增加了消息的可靠性,但也不好说能保证消息100%的可靠。万一真的MQ通知失败该怎么办呢?有没有其它补偿方案,能够确保订单的支付状态一致呢?

其实思想很简单:既然MQ通知不一定发送到交易服务,那么交易服务就必须自己主动去查询支付状态。这样即便支付服务的MQ通知失败,我们依然能通过主动查询来保证订单状态的一致。

流程如下:

画板2

图中黄色线圈起来的部分就是MQ通知失败后的补偿处理方案,由交易服务自己主动去查询支付状态。

不过需要注意的是,交易服务并不知道用户会在什么时候支付,如果查询的时机不正确(比如查询的时候用户正在支付中),可能查询到的支付状态也不正确。

那么问题来了,我们到底该在什么时间主动查询支付状态呢?

  1. 由用户在页面中主动点击“已完成支付”按钮时去查询最新的支付结果。

配图16

  1. 通过定时任务去定义查询支付结果,比如:对于30分钟内的订单状态为未支付状态时去主动查询支付结果,超过30分钟用户未支付将自动取消订单。

可以每隔一段时间查询一次,并判断支付状态。如果发现订单已经支付,则立刻更新订单状态为已支付即可。

定时任务大家之前学习过,具体的实现这里就不再赘述了。

至此,消息可靠性的问题已经解决了。

综上,支付服务与交易服务之间的订单状态一致性是如何保证的?

  • 首先,支付服务在用户支付成功以后利用MQ消息通知交易服务,完成订单状态同步。
  • 其次,为了保证MQ消息的可靠性,我们采用了生产者确认机制、消费者确认、消费者失败重试等策略,确保消息投递的可靠性
  • 最后,我们提供了用户主动查询支付结果的入口 ,并且还在交易服务设置了定时任务,定期查询订单支付状态。这样即便MQ通知失败,还可以利用定时任务作为兜底方案,确保订单支付状态的最终一致性。

1.4.5. 面试题 ​

如何保证消息的可靠性?可以百分百保证MQ的消息可靠性吗?

如何保证MQ幂等性?或 如何防止消息重复消费?

2.延迟消息 ​

2.1. 技术方案 ​

在电商的支付业务中,对于一些库存有限的商品,为了更好的用户体验,通常都会在用户下单时立刻扣减商品库存。例如电影院购票、高铁购票,下单后就会锁定座位资源,其他人无法重复购买。

但是这样就存在一个问题,假如用户下单后一直不付款,就会一直占有库存资源,导致其他客户无法正常交易,最终导致商户利益受损!

因此,电商中通常的做法就是:对于超过一定时间未支付的订单会自动取消订单并释放占用的库存。

例如,订单支付超时时间为30分钟,则我们应该在用户下单后的第30分钟检查订单支付状态,如果发现未支付,应该立刻取消订单,释放库存。

但问题来了:如何才能准确的实现在下单后第30分钟去检查支付状态呢?

像这种在一段时间以后才执行的任务,我们称之为延迟任务,而要实现延迟任务,最简单的方案就是利用MQ的延迟消息了。

在RabbitMQ中实现延迟消息也有两种方案:

2.1.1 死信交换机+TTL ​

该方案利用死信交换机实现,死信交换机(Dead Letter Exchange, DLX)是一种处理消息队列中无法被消费的消息的方式。当消息因为某些原因(如消费者拒绝消费消息、消息过期等)而无法被正常消费时,这些消息就会成为“死信”。

结合 TTL(Time To Live),你可以设定消息在队列中的存活时间,当消息在队列中停留的时间超过了设定的 TTL 后,消息就会成为死信。当消息变成死信时,会被发送到死信交换机中去,通过死信交换机转发到指定的队列,由应用程序去消费,进一步处理这些消息。

假设需求是下单后30分钟对未支付订单自动取消订单,通过下图举例说明:

配图17

流程如下:

  1. 创建订单成功向“ttl.fanout”交换机发送消息,消息内容记录下单的订单信息,消息的TTL为30分钟
  2. 通过"ttl.fanout"将消息转发到"ttl.queue"队列
  3. 由于该队列并没有消费程序去监听,当到达30分钟时该消息变为死信,自动发给死信交换机“hmall.direct”,由列信交换机将消息转发到"direct.queue1"队列。
  4. 应用监听“direct.queue1”队列,收到未支付超时订单,执行取消订单并回滚库存。

此方案的问题是:

当一个消息的TTL到期以后不一定会被移除或投递到死信交换机,而是在消息恰好处于队首时才会被处理。

当队列中消息堆积很多的时候,过期消息可能不会被按时处理,因此你设置的TTL时间不一定准确。

2.1.2 延迟消息插件 ​

RabbitMQ社区提供了一个延迟消息插件来实现相同的效果。

此方案实现过程非常简单,安装延迟消息插件,发送消息时指定延迟时间,到达延迟时间消息被消费。

延迟消息插件内部会维护一个本地数据库表,同时使用Elang Timers功能实现计时。如果消息的延迟时间设置较长,可能会导致堆积的延迟消息非常多,会带来较大的CPU开销,同时延迟消息的时间会存在误差。

因此,不建议设置延迟时间过长的延迟消息。

对上述两个方案比较,实现延迟消息推荐使用延迟消息插件实现。

2.2.使用延迟消息插件 ​

官方文档说明:

2.2.1下载 ​

插件下载地址:

由于我们安装的MQ是3.8版本,因此这里下载3.8.17版本:

配图18

当然,也可以直接使用课前资料提供好的插件:

配图19

2.2.2 安装 ​

因为我们是基于Docker安装,所以需要先查看RabbitMQ的插件目录对应的数据卷。

Shell
docker volume inspect mq-plugins
1

结果如下:

JSON
[
    {
        "CreatedAt": "2024-06-19T09:22:59+08:00",
        "Driver": "local",
        "Labels": null,
        "Mountpoint": "/var/lib/docker/volumes/mq-plugins/_data",
        "Name": "mq-plugins",
        "Options": null,
        "Scope": "local"
    }
]
1
2
3
4
5
6
7
8
9
10
11

插件目录被挂载到了/var/lib/docker/volumes/mq-plugins/_data这个目录,我们上传插件到该目录下。

接下来执行命令,安装插件:

Shell
docker exec -it mq rabbitmq-plugins enable rabbitmq_delayed_message_exchange
1

运行结果如下:

配图20

2.2.3 声明延迟交换机 ​

控制台方式:

配图21

创建delay.queue队列,并绑定到延迟交换机

配图22

基于注解方式:

Java
@RabbitListener(bindings = @QueueBinding(
        value = @Queue(name = "delay.queue", durable = "true"),
        exchange = @Exchange(name = "delay.direct", delayed = "true",type = ExchangeTypes.DIRECT, durable = "true"),
        key = "delay"
))
public void listenDelayMessage(String msg){
    log.info("接收到delay.queue的延迟消息:{}", msg);
}
1
2
3
4
5
6
7
8

2.2.4 发送延迟消息 ​

发送消息时,必须通过x-delay属性设定延迟时间:

Java
@Test
void testPublisherDelayMessage() {
    // 1.创建消息
    String message = "hello, delayed message";
    // 2.发送消息,利用消息后置处理器添加消息头
    rabbitTemplate.convertAndSend("delay.direct", "delay", message, new MessagePostProcessor() {
        @Override
        public Message postProcessMessage(Message message) throws AmqpException {
            // 添加延迟消息属性
            message.getMessageProperties().setDelay(5000);
            log.info("发送消息"+new String(message.getBody(), StandardCharsets.UTF_8)+ LocalDateTime.now());
            return message;
        }
    });
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15

2.2.5 测试 ​

测试流程:

  1. 启动消费端监听程序。

在consumer模块添加监听方法。

Java
@RabbitListener(queues = "delay.queue")
public void listenDelayQueue(String msg) throws InterruptedException {
    System.out.println("接收延迟消息:【" + msg + "】" + LocalTime.now());
}
1
2
3
4
  1. 执行发送程序
  2. 观察发送程序日志中发送的消息的时间,观察监听程序的日志中接收消息时间的,判断是否是按延迟时间发送消息。

注意:发送延迟消息会触发ReturnCallback,这是因为延迟消息并没有发送到队列,而是在延迟时间到达才发送到队列,所以会触发ReturnCallback,需要在ReturnCallback方法中判断,如果是延迟消息不要向失败消息表插入记录。

2.3.自动取消超时未支付订单(作业) ​

2.3.1. 需求分析 ​

接下来,我们就在交易服务中利用延迟消息实现订单超时取消功能,假设订单未支付超时时间为30分钟,下单后30内不支付将自动取消订单,其大概思路如下:

配图23

理论上我们应该在下单时发送一条延迟消息,延迟时间为30分钟。这样就可以在接收到消息时检验订单支付状态,关闭未支付订单。

不过大家设想,极端情况如果用户在下单后正好30分钟去支付了,此时延迟消息发送到交易服务,交易服务收到延迟消息后取消订单,用户支付成功,问题出现。

我们可以把延迟消息的时间定为35分钟,当下单30分钟后用户无法支付,35分钟后自动取消订单避免取消订单后用户支付成功的问题。

交互流程如下:

画板3

2.3.2 基础配置 ​

我们在常量类中配置交换机、队列、RoutingKey等常量,内容如下:

Java
package com.hmall.common.constants;

public interface MqConstants{
    String DELAY_EXCHANGE_NAME = "trade.delay.direct";
    String DELAY_ORDER_QUEUE_NAME = "trade.delay.order.queue";
    String DELAY_ORDER_KEY = "delay.order.query";
}
1
2
3
4
5
6
7

在trade-service模块的pom.xml中引入amqp的依赖:

XML
  
  
      org.springframework.boot
      spring-boot-starter-amqp
1
2
3
4
5

在trade-service的application.yaml中添加MQ的配置:

YAML
spring:
  rabbitmq:
    host: 192.168.101.68
    port: 5672
    virtual-host: /hmall
    username: hmall
    password: 123
1
2
3
4
5
6
7

2.3.3. 发送延迟消息 ​

接下来,我们改造下单业务,在下单完成后,发送延迟消息。

修改trade-service模块的com.hmall.trade.service.impl.OrderServiceImpl类的createOrder方法,添加消息发送的代码:

JavaScript
    @Override
//    @Transactional
    @GlobalTransactional
    public Long createOrder(OrderFormDTO orderFormDTO) {
    ....
    ....
    //todo 发送延迟消息
    
    return order.getId();
}
1
2
3
4
5
6
7
8
9
10

这里延迟消息的时间应该是35分钟,不过我们为了测试方便暂时设置为10秒。

2.3.4. 接收延迟消息 ​

接下来,我们在trader-service编写一个监听器,监听延迟消息,查询订单支付状态:

配图24

代码如下:

Java
package com.hmall.trade.listener;

import com.hmall.api.client.PayClient;
import com.hmall.api.dto.PayOrderDTO;
import com.hmall.common.constants.MqConstants;
import com.hmall.trade.domain.po.Order;
import com.hmall.trade.service.IOrderService;
import lombok.RequiredArgsConstructor;
import org.springframework.amqp.rabbit.annotation.Exchange;
import org.springframework.amqp.rabbit.annotation.Queue;
import org.springframework.amqp.rabbit.annotation.QueueBinding;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
@RequiredArgsConstructor
public class OrderDelayMessageListener {

    private final IOrderService orderService;

        @RabbitListener(bindings = @QueueBinding(
        value = @Queue(name = MqConstants.DELAY_ORDER_QUEUE_NAME),
        exchange = @Exchange(name = MqConstants.DELAY_EXCHANGE_NAME, delayed = "true"),
        key = MqConstants.DELAY_ORDER_KEY
        ))
        public void listenOrderDelayMessage(Long orderId) {
            // 1.查询订单
            Order order = orderService.getById(orderId);
            // 2.检测订单状态,判断是否已支付
            if (order == null || order.getStatus() != 1) {
                // 订单不存在或者已经支付
                return;
            }
            // 3.todo 未支付,需要查询支付流水状态

            // 4.判断是否支付

            // 4.1.已支付,标记订单状态为已支付
            orderService.markOrderPaySuccess(orderId);

            // TODO 4.2.未支付,取消订单,回滚库存
                
        }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44

逻辑如下:

1、收到消息表示支付成功,查询支付系统最新的支付状态

2、如果已支付则标记订单为已支付

3、如果未支付则取消订单

取消订单需要完成两件事情:

  • 将订单状态修改为已关闭
  • 将订单中已经扣除的库存重新加回来

大家在IOrderService接口中定义cancelOrder方法,此方法为取消订单的service方法。

Java
void cancelOrder(Long orderId);
1

并且在OrderServiceImpl中实现该方法,实现过程中要注意业务幂等性判断。

2.3.5. 编写查询支付状态接口 ​

由于MQ消息处理时需要查询当前最新的支付状态,因此我们要在pay-service模块定义查询支付状态的接口

在pay-service模块的PayController中实现该接口:

Java
@ApiOperation("根据id查询支付单")
@GetMapping("/biz/{id}")
public PayOrderDTO queryPayOrderByBizOrderNo(@PathVariable("id") Long id){
    PayOrder payOrder = payOrderService.lambdaQuery().eq(PayOrder::getBizOrderNo, id).one();
    return BeanUtils.copyBean(payOrder, PayOrderDTO.class);
}
1
2
3
4
5
6

此接口会在api工程定义,所以将PayOrderDTO 在api工程定义,PayOrderDTO 的结果同PayOrder ,为防止返回数据有变这里定义为DTO类型,编写时可直接将PayOrder 拷贝至PayOrderDTO 中并进行微改。

PayOrderDTO代码如下:

Java
package com.hmall.api.pay.dto;

import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;

import java.time.LocalDateTime;

/**
 * 
 * 支付订单
 * 
 */
@Data
@ApiModel(description = "支付单数据传输实体")
public class PayOrderDTO {
    @ApiModelProperty("id")
    private Long id;
    @ApiModelProperty("业务订单号")
    private Long bizOrderNo;
    @ApiModelProperty("支付单号")
    private Long payOrderNo;
    @ApiModelProperty("支付用户id")
    private Long bizUserId;
    @ApiModelProperty("支付渠道编码")
    private String payChannelCode;
    @ApiModelProperty("支付金额,单位分")
    private Integer amount;
    @ApiModelProperty("付类型,1:h5,2:小程序,3:公众号,4:扫码,5:余额支付")
    private Integer payType;
    @ApiModelProperty("付状态,0:待提交,1:待支付,2:支付超时或取消,3:支付成功")
    private Integer status;
    @ApiModelProperty("拓展字段,用于传递不同渠道单独处理的字段")
    private String expandJson;
    @ApiModelProperty("第三方返回业务码")
    private String resultCode;
    @ApiModelProperty("第三方返回提示信息")
    private String resultMsg;
    @ApiModelProperty("支付成功时间")
    private LocalDateTime paySuccessTime;
    @ApiModelProperty("支付超时时间")
    private LocalDateTime payOverTime;
    @ApiModelProperty("支付二维码链接")
    private String qrCodeUrl;
    @ApiModelProperty("创建时间")
    private LocalDateTime createTime;
    @ApiModelProperty("更新时间")
    private LocalDateTime updateTime;
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49

接下来定义接口对应的FeignClient.

PayClient代码如下:

Java
package com.hmall.api.client;

import com.hmall.api.client.fallback.PayClientFallback;
import com.hmall.api.dto.PayOrderDTO;
import org.springframework.cloud.openfeign.FeignClient;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;

@FeignClient(value = "pay-service", fallbackFactory = PayClientFallbackFactory .class)
public interface PayClient {
    /**
     * 根据交易订单id查询支付单
     * @param id 业务订单id
     * @return 支付单信息
     */
    @GetMapping("/pay-orders/biz/{id}")
    PayOrderDTO queryPayOrderByBizOrderNo(@PathVariable("id") Long id);
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18

PayClientFallbackFactory 代码如下:

Java
package com.hmall.api.client.fallback;

import com.hmall.api.client.PayClient;
import com.hmall.api.dto.PayOrderDTO;
import lombok.extern.slf4j.Slf4j;
import org.springframework.cloud.openfeign.FallbackFactory;

@Slf4j
public class PayClientFallbackFactory implements FallbackFactory {
    @Override
    public PayClient create(Throwable cause) {
        return new PayClient() {
            @Override
            public PayOrderDTO queryPayOrderByBizOrderNo(Long id) {
                return null;
            }
        };
    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19

最终在hm-api模块定义三个类:

配图25

说明:

  • PayOrderDTO:支付单的数据传输实体
  • PayClient:支付系统的Feign客户端
  • PayClientFallback:支付系统的fallback逻辑

3.MQ集群 ​

生产环境通常会搭建MQ集群,通常运维人员负责,如果工作中让Java工程师负责我们只需要参考部署文档进行部署即可。

参考:RabbitMQ集群部署

通过本节对MQ集群的演示理解集群方式不仅可以提高性能还可以实现高可用。

虚拟机中mq1、mq2、mq3镜像模式集群已搭建好。

首先通过文档理解非镜像模式的问题。

再演示镜像集群的优势。

演示仲裁队列。

4 Elasticsearch快速入门 ​

4.1. 认识Elasticsearch ​

4.1.1 搜索需求 ​

黑马商城作为一个电商项目,商品的搜索肯定是访问频率最高的页面之一。目前搜索功能是基于数据库的模糊搜索来实现的,存在很多问题。

查询效率较低 ​

由于数据库模糊查询不走索引,在数据量较大的时候,查询性能很差。黑马商城的商品表中仅仅有不到9万条数据,基于数据库查询时,搜索接口的表现如图:

配图26

改为基于搜索引擎后,查询表现如下:

配图27

需要注意的是,数据库模糊查询随着表数据量的增多,查询性能的下降会非常明显,而搜索引擎的性能则不会随着数据增多而下降太多。目前仅10万不到的数据量差距就如此明显,如果数据量达到百万、千万、甚至上亿级别,这个性能差距会非常夸张。

功能单一 ​

数据库的模糊搜索功能单一,匹配条件非常苛刻,必须恰好包含用户搜索的关键字。而在搜索引擎中,用户输入出现个别错字,或者用拼音搜索、同义词搜索都能正确匹配到数据。

综上,在面临海量数据的搜索,或者有一些复杂搜索需求的时候,推荐使用全文搜索技术来实现搜索功能。

AI:开源全文搜索引擎,按使用用户量降序排序

根据社区支持、知名度和普遍的应用情况,我们可以做出一个大致的估计。以下是几个较为知名的开源全文搜索引擎,按照它们的流行度和使用广泛程度进行排序:

  1. Elasticsearch

    • 概述:Elasticsearch 是一个分布式的、RESTful 接口的搜索引擎和分析引擎,基于 Lucene 构建。
    • 使用情况:被广泛用于日志分析、全文搜索、实时分析等领域,是 Elastic Stack 的核心组件之一。
    • 社区支持:拥有庞大的社区和丰富的文档,以及大量的第三方插件和工具。
    • 官网:https://www.elastic.co/elasticsearch
  2. Apache Solr

    • 概述:Solr 是一个高性能的企业级搜索平台,同样基于 Lucene 构建。
    • 使用情况:广泛应用于电子商务、新闻媒体、社交网络等多个领域。
    • 社区支持:拥有成熟的社区和文档,支持多种语言和框架。
    • 官网:https://solr.apache.org/
  3. Apache Lucene

    • 概述:Lucene 是一个高性能的全文检索库,是许多搜索引擎的基础。
    • 使用情况:由于其高度可定制性,Lucene 被广泛应用于构建定制化的搜索解决方案。
    • 社区支持:拥有庞大的 Java 开发者社区,支持多种编程语言。
    • 官网:http://lucene.apache.org/
  4. MeiliSearch

    • 概述:MeiliSearch 是一个高度可配置、易于使用的搜索引擎。
    • 使用情况:适用于需要快速部署的项目,特别是在开发阶段。
    • 社区支持:拥有活跃的社区和良好的文档。
    • 官网:https://meilisearch.com/
  5. Zinc

    • 概述:Zinc 是一个轻量级的全文搜索引擎,用 Go 语言编写,旨在作为 Elasticsearch 的轻量级替代方案。
    • 使用情况:适用于资源有限的环境,如嵌入式系统或小型项目。
    • 社区支持:相对较小的社区,但持续增长。
    • 官网:https://github.com/justwatchcom/zinc
  6. CloriSearch

    • 概述:CloriSearch 是一个轻量级的全文搜索引擎,用 Rust 语言编写。
    • 使用情况:适用于需要高性能和稳定性的项目。
    • 社区支持:社区正在成长中,但提供了一个简洁且强大的接口。
    • 官网:https://gitcode.net/shpilu/cloriSearch

排名第一的就是我们今天要学习的Elasticsearch.

Elasticsearch是一款非常强大的开源搜索引擎,支持的功能非常多,例如:

配图28代码搜索

配图29商品搜索

配图30解决方案搜索

配图31地图搜索

4.1.2 倒排索引 ​

Elasticsearch之所以有如此高性能的搜索表现,正是得益于底层的倒排索引技术。那么什么是倒排索引呢?

倒排索引的概念是基于正向索引而言的。

4.1.2.1 正向索引 ​

我们先来回顾一下正向索引。

例如有一张名为tb_goods的表:

idtitleprice
1小米手机3499
2华为手机4999
3华为小米充电器49
4小米手环49
.........

其中的id字段已经创建了索引,由于索引底层采用了B+树结构,因此我们根据id搜索的速度会非常快。但是其他字段例如title,只在叶子节点上存在。

因此要根据title搜索的时候只能遍历树中的每一个叶子节点,判断title数据是否符合要求。

比如用户的SQL语句为:

SQL
select * from tb_goods where title like '%手机%';
1

那搜索的大概流程如图:

配图32

说明:

  • 1)检查到搜索条件为like '%手机%',需要找到title中包含手机的数据
  • 2)逐条遍历每行数据(每个叶子节点),比如第1次拿到id为1的数据
  • 3)判断数据中的title字段值是否符合条件
  • 4)如果符合则放入结果集,不符合则丢弃
  • 5)回到步骤1

综上,根据id精确匹配时,可以走索引,查询效率较高。而当搜索条件为模糊匹配时,由于索引无法生效,导致从索引查询退化为全表扫描,效率很差。

因此,正向索引适合于根据索引字段的精确搜索,不适合基于部分词条的模糊匹配。

而倒排索引恰好解决的就是根据部分词条模糊匹配的问题。

4.1.2.2 倒排索引 ​

倒排索引中有两个非常重要的概念:

  • 文档(Document):用来搜索的数据,其中的每一条数据就是一个文档。例如一个网页、一个商品信息
  • 词条(Term):对文档数据或用户搜索数据,利用某种算法分词,得到的具备含义的词语就是词条。例如:我是中国人,就可以分为:我、是、中国人、中国、国人这样的几个词条

创建倒排索引是对正向索引的一种特殊处理和应用,流程如下:

  • 将每一个文档的数据利用分词算法根据语义拆分,得到一个个词条
  • 倒排索引记录每个词条对应的文档id

此时形成的这张以词条为索引的表,就是倒排索引表,两者对比如下:

正向索引

id(索引)titleprice
1小米手机3499
2华为手机4999
3华为小米充电器49
4小米手环49
.........

倒排索引

词条(索引)文档id
小米1,3,4
手机1,2
华为2,3
充电器3
手环4

倒排索引的搜索流程如下(以搜索"华为手机"为例),如图:

配图33

流程描述:

1)用户输入条件"华为手机"进行搜索。

2)对用户输入条件分词,得到词条:华为、手机。

3)拿着词条在倒排索引中查找(由于词条有索引,查询效率很高),即可得到包含词条的文档id:1、2、3。

4)拿着文档id到正向索引中查找具体文档即可(由于id也有索引,查询效率也很高)。

虽然要先查询倒排索引,再查询倒排索引,但是无论是词条、还是文档id都建立了索引,查询速度非常快!无需全表扫描。

4.1.2.3 面试题 ​

说下倒排索引?为什么叫倒排索引?

项目为什么要用Elasticsearch?数据很多吗?

4.2. 安装Elasticsearch ​

4.2.1 介绍Elasticsearch ​

Elasticsearch是由elastic公司开发的一套搜索引擎技术,Elasticsearch的官方网站如下:

配图34

它是elastic技术栈中的一部分,整套技术栈被称为ELK,经常用来做日志收集、系统监控和状态分析等等。

AI:ELK是干什么的

ELK 是指一套流行的技术栈,用于构建高级的日志收集、分析和展示系统。ELK 是三个组件首字母的缩写,它们分别是:

  1. Elasticsearch:这是一个分布式的、RESTful接口的搜索和分析引擎。它能够存储大量的数据,并且提供近乎实时的搜索能力。在ELK栈中,Elasticsearch作为存储和搜索日志数据的核心。
  2. Logstash:Logstash 是一个服务器端的工具,用于接收来自各种数据源的日志数据,它可以解析、转换这些数据,并将其发送到Elasticsearch中存储。Logstash支持多种插件,可以轻松地从不同的数据源收集数据。
  3. Kibana:Kibana 是一个基于Web的用户界面,用于对存储在Elasticsearch中的数据进行可视化。它允许用户创建复杂的查询来分析数据,并以图表、表格等形式展示数据。

近年来,随着Elastic生态的发展,Beats(比如Filebeat)这类轻量级的日志收集器也被广泛采用,它们通常部署在多个服务器上,负责将日志数据发送给Logstash或者直接送入Elasticsearch。尽管Beats不是传统意义上的ELK栈的一部分,但是它们经常被一起提及,因为它们简化了日志数据的收集过程。

总的来说,ELK栈提供了一个全面的日志管理系统,它能够帮助组织捕获、存储、分析和展示大量的日志数据。这套系统常被用来监测应用程序性能、跟踪用户行为、进行网络安全分析等多种用途。

配图35

整套技术栈的核心就是用来存储、搜索、计算的Elasticsearch,因此我们接下来学习的核心也是Elasticsearch。

我们要安装的内容包含2部分:

  • elasticsearch:存储、搜索和运算
  • kibana:图形化展示控制台

4.2.2 安装Elasticsearch ​

我们当前使用的Spring Boot2.7.X版本默认使用的是Elasitcsearch7.17.x,本课程基于7.17.7版本学习。

通过下面的Docker命令即可安装单机版本的elasticsearch:

拉取镜像

docker pull elasticsearch:7.17.7

由于镜像较大也可将课程资料中“es安装”目录下的elasticsearch.7.17.7.tar上传到虚拟机,然后导入docker镜像,执行下边的命令:

docker load -i elasticsearch.7.17.7.tar

创建文件夹:

mkdir -p /data/soft/es7.17.7/xzb

在/data/soft/es7.17.7/xzb下创建data目录并且修改权限为777

Java
mkdir data
chmod 777 data
1
2

将课程资料下的"ES安装"目录中的 es.zip上传到/data/soft/es7.17.7/xzb下,并进行解压

Java
 unzip es.zip
1

解压成功如下图:

配图36

创建容器

Java
docker run -d \
--name elasticsearch7.17.7 \
--restart always \
-p 9200:9200 \
-p 9300:9300 \
-e "discovery.type=single-node" \
-e "ES_JAVA_OPTS=-Xms512m -Xmx512m" \
-v /data/soft/es7.17.7/xzb/data:/usr/share/elasticsearch/data \
-v /data/soft/es7.17.7/xzb/plugins:/usr/share/elasticsearch/plugins \
-v /data/soft/es7.17.7/xzb/config:/usr/share/elasticsearch/config \
elasticsearch:7.17.7
1
2
3
4
5
6
7
8
9
10
11

安装完成后,访问9200端口(http://192.168.101.68:9200/),即可看到响应的Elasticsearch服务的基本信息:

JavaScript
{
  "name" : "4251f98ff357",
  "cluster_name" : "docker-cluster",
  "cluster_uuid" : "aB_5c-y4St-NU-MFHxiVvg",
  "version" : {
    "number" : "7.17.7",
    "build_flavor" : "default",
    "build_type" : "docker",
    "build_hash" : "78dcaaa8cee33438b91eca7f5c7f56a70fec9e80",
    "build_date" : "2022-10-17T15:29:54.167373105Z",
    "build_snapshot" : false,
    "lucene_version" : "8.11.1",
    "minimum_wire_compatibility_version" : "6.8.0",
    "minimum_index_compatibility_version" : "6.0.0-beta1"
  },
  "tagline" : "You Know, for Search"
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17

4.2.3 安装Kibana ​

通过下面的Docker命令,即可部署Kibana:

拉取镜像

docker pull kibana:7.17.7

由于镜像较大也可将课程资料中“es安装”目录下的kibana.7.17.7.tar 上传到虚拟机,然后导入docker镜像,执行下边的命令:

docker load -i kibana.7.17.7.tar

创建容器:

注意修改es的地址

Java
docker run --name kibana7.17.7 \
-e ELASTICSEARCH_HOSTS=http://192.168.101.68:9200 \
-p 5601:5601 \
-d kibana:7.17.7
1
2
3
4

下边启动容器,先保证Elasticsearch启动成功。

启动kibana容器成功,在浏览器输入地址访问:http://192.168.101.68:5601

配图37

4.2.4 小结 ​

安装Elasticsearch和Kibana需要注意:Elasticsearch和Kibana的版本需要保持一致。

我们项目用的版本是7.17.7。

ELK是干什么的?包括哪些中间件?

ELK用于构建日志收集分析系统,包括:

  • Elasticsearch:用于数据存储、计算和搜索
  • Logstash/Beats:用于数据收集
  • Kibana:用于数据可视化

通过Logstash将应用程序的日志采集到Elasticsearch中,通过Elasticsearch对日志进行分析,通过Kibana展示查询日志,展示分析的结果。

4.3.基础概念 ​

elasticsearch中有很多独有的概念,与mysql中略有差别,但也有相似之处。

4.3.1 文档和字段 ​

elasticsearch是面向文档(Document)存储的,可以是数据库中的一条商品数据,一个订单信息。文档数据会被序列化为json格式后存储在elasticsearch中:

配图38

JSON
{
    "id": 1,
    "title": "小米手机",
    "price": 3499
}
{
    "id": 2,
    "title": "华为手机",
    "price": 4999
}
{
    "id": 3,
    "title": "华为小米充电器",
    "price": 49
}
{
    "id": 4,
    "title": "小米手环",
    "price": 299
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21

因此,原本数据库中的一行数据就是ES中的一个JSON文档;而数据库中每行数据都包含很多列,这些列就转换为JSON文档中的字段(Field)。

4.3.2 索引和映射 ​

随着业务发展,需要在es中存储的文档也会越来越多,比如有商品的文档、用户的文档、订单文档等等:

配图39

所有文档都散乱存放显然非常混乱,也不方便管理,因此,我们要将相同类型的文档集中在一起管理,称为索引(Index)。例如:

商品索引

JSON
{
    "id": 1,
    "title": "小米手机",
    "price": 3499
}

{
    "id": 2,
    "title": "华为手机",
    "price": 4999
}

{
    "id": 3,
    "title": "三星手机",
    "price": 3999
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17

用户索引

JSON
{
    "id": 101,
    "name": "张三",
    "age": 21
}

{
    "id": 102,
    "name": "李四",
    "age": 24
}

{
    "id": 103,
    "name": "麻子",
    "age": 18
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17

订单索引

JSON
{
    "id": 10,
    "userId": 101,
    "goodsId": 1,
    "totalFee": 294
}

{
    "id": 11,
    "userId": 102,
    "goodsId": 2,
    "totalFee": 328
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
  • 所有用户文档,就可以组织在一起,称为用户的索引;
  • 所有商品的文档,可以组织在一起,称为商品的索引;
  • 所有订单的文档,可以组织在一起,称为订单的索引;

索引就类似数据库表,MySQL中我们会先创建表结构再向表中插入数据,同样,ES中的索引也有结构,那就是映射(mapping),在 Elasticsearch 中,映射(mapping)定义了索引(index)中文档(document)的结构和字段(field)的数据类型及属性,映射类似于关系数据库中的表结构定义,它告诉 Elasticsearch 如何解析、存储和索引数据。

4.3.3 总结 ​

我们对mysql与elasticsearch的概念做一下对比:

MySQLElasticsearch说明
TableIndex索引(index),就是文档的集合,类似数据库的表(table)
RowDocument文档(Document),就是一条条的数据,类似数据库中的行(Row),文档都是JSON格式
ColumnField字段(Field),就是JSON文档中的字段,类似数据库中的列(Column)
SchemaMappingMapping(映射)是索引中文档的约束,例如字段类型约束。类似数据库的表结构(Schema)
SQLDSLDSL是elasticsearch提供的JSON风格的请求语句,用来操作elasticsearch,实现CRUD

如图:

配图40

那是不是说,我们学习了elasticsearch就不再需要mysql了呢?

并不是如此,两者各自有自己的擅长之处:

  • Mysql:擅长事务类型操作,可以确保数据的安全和一致性
  • Elasticsearch:擅长海量数据的搜索、分析、计算

因此在企业中,往往是两者结合使用:

  • 对安全性要求较高的写操作,使用mysql实现
  • 对查询性能要求较高的搜索需求,使用elasticsearch实现
  • 两者再基于某种方式,实现数据的同步,保证一致性

配图41

4.4. 快速入门 ​

4.4.1 创建索引 ​

根据前边对倒排索引的理解,倒排索引就是根据词找文档,词就是索引,所以要想完成搜索功能开发第一步就是要创建索引,有了索引就可以搜索了。

打开Kibana,进入DevTools,如下图:

配图42

进入DevTools

配图43

AI: elasticsearch快速入门

执行下边的命令向ES添加文档,如果my_index索引不存在会自动创建:

JavaScript
POST /my_index/_doc/1
{
  "title": "Elasticsearch: cool and easy",
  "content": "This is a test document"
}
1
2
3
4
5
6

Elasticsearch提供RESTful接口供创建索引、修改索引、删除索引等操作。

请求路径:/my_index/_doc/1

请求内容:json结构

整体路径表示一个文档的地址。

my_index:表示索引名称,相当于MySQL的表名,如果没有会自动创建。

_doc:索引类型(type), 在Elasticsearch 7.x 版本之前一个索引中的文档可以归属不同的类型,这样非常不好理解,从7.x 及之后 统一使用 _doc 作为索引的类型,也就是不存在类型这个概念了,固定写为_doc即可。

1: 是文档的唯一标识符(ID)。在 Elasticsearch 中,每个文档都有一个唯一的 ID,相当于MySQL中一个表的主键值。

Elasticsearch会对title、content两个字段的内容进行分词,每个词条关联1号文档。

"Elasticsearch: cool and easy" 分词为:Elasticsearch、cool、and、easy,默认分词器按空格分词。

"This is a test document" 分词为:this、is、a、test、document

执行结果如下图:

配图44

结果显示:索引名称为my_index

successful: 插入成功1个文档。

4.4.2 查询文档 ​

根据id查询文档:GET /my_index/_doc/1

my_index:索引名

_doc: 固定

1: 文档的id

查询结果

配图45

4.4.3 搜索文档 ​

下边进行搜索:

执行下边的命令:

JavaScript
GET /my_index/_search
{
  "query": {
    "match": {
      "content": "test"
    }
  }
}
1
2
3
4
5
6
7
8

说明:

content:my_index索引中的字段名。

"test": 搜索的关键字。

搜索流程:

  1. 对搜索的关键字进行分词
  2. 拿的词去索引中搜索,最终找到匹配分词的文档。

执行结果:

配图46

动手实验,下边的搜索可以搜出结果吗?

JavaScript
GET /my_index/_search
{
  "query": {
    "match": {
      "content": "test java"
    }
  }
}
1
2
3
4
5
6
7
8

4.4.4 删除文档 ​

执行下边的语句删除文档:

JavaScript
DELETE /my_index/_doc/1
1

如下图:

配图47

删除了1号文档,此时再去搜索1号文档的内容还可以找到吗?

执行下边的搜索再试试:

JavaScript
GET /my_index/_search
{
  "query": {
    "match": {
      "content": "test java"
    }
  }
}
1
2
3
4
5
6
7
8

4.4.5 总结 ​

通过Elasticsearch提供的RESTful接口操作索引:

1、添加索引

将文档信息提交给Elasticsearch,指定索引名称及文档内容,它对文档内容进行分词、存储。

索引的结构是倒排索引表。

2、搜索

根据文档id查询文档

指定文档的字段及搜索关键字进行搜索。

搜索过程会先将关键字进行分词,再拿词去索引中查询。

基本概念:

相关概念如下:

  1. 索引 (index):

    • Elasticsearch 中的数据被组织成索引,每个索引都有一个唯一的名称。
    • 一个索引可以包含多个文档。
  2. 文档 (document):

    • 文档是 Elasticsearch 中的基本单位,每个文档都是一个 JSON 对象。
    • 文档包含一个或多个字段。
  3. 字段 (field):

    • 字段是文档中的基本单元,用于存储数据。
    • 每个字段都有一个数据类型,例如 text、keyword、integer、float 等。

4.5. IK分词器 ​

4.5.1 认识分词器 ​

Elasticsearch的底层是倒排索引,倒排索引中的词条来源于对文档内容的分词,分词器正是负责对文档内容进行分词。

Elasticsearch 的分词器(analyzers)是用于处理文本数据的关键组件。它们负责将原始文本分解成一系列词条(tokens),并对这些词条进行规范化(normalization),以便进行索引和搜索。分词器由两个主要部分组成:tokenizer(分词) 和 filter(过滤)。Tokenizer 负责将文本分割成词条,而 filter 则对这些词条进行额外的处理,比如转换为小写、去除停用词等。

Elasticsearch 提供了多种内置分词器:

Standard Analyzer:标准分词器(standard),这是默认的分词器,用于大多数情况。它会移除标点符号,并将文本转换为小写。

Stop Analyzer:停用词分词器(stop),除了执行标准分词器的操作之外,还会过滤掉一些常见的英文停用词(stop words)。

Simple Analyzer:简单分词器(simple),会移除常见的 HTML 标签,并且会将所有字母转换为小写。

Whitespace Analyzer:空白字符分词器(whitespace),仅根据空白字符分割文本。

下边可以测试标准分词器的分词效果:

标准分词器是根据空格进行分词。

"analyzer"指定分词器名称"standard"

"text":分词内容

JavaScript

POST _analyze
{
  "analyzer": "standard",
  "text": "The quick brown fox jumps over the lazy dog."
}
1
2
3
4
5
6
7
8

效果:

JavaScript

{
  "tokens" : [
    {
      "token" : "the",
      "start_offset" : 0,
      "end_offset" : 3,
      "type" : "",
      "position" : 0
    },
    {
      "token" : "quick",
      "start_offset" : 4,
      "end_offset" : 9,
      "type" : "",
      "position" : 1
    },
    {
      "token" : "brown",
      "start_offset" : 10,
      "end_offset" : 15,
      "type" : "",
      "position" : 2
    },
    {
      "token" : "fox",
      "start_offset" : 16,
      "end_offset" : 19,
      "type" : "",
      "position" : 3
    },
    {
      "token" : "jumps",
      "start_offset" : 20,
      "end_offset" : 25,
      "type" : "",
      "position" : 4
    },
    {
      "token" : "over",
      "start_offset" : 26,
      "end_offset" : 30,
      "type" : "",
      "position" : 5
    },
    {
      "token" : "the",
      "start_offset" : 31,
      "end_offset" : 34,
      "type" : "",
      "position" : 6
    },
    {
      "token" : "lazy",
      "start_offset" : 35,
      "end_offset" : 39,
      "type" : "",
      "position" : 7
    },
    {
      "token" : "dog",
      "start_offset" : 40,
      "end_offset" : 43,
      "type" : "",
      "position" : 8
    }
  ]
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69

再测试停用词分词器:

在 Elasticsearch 中,停用词(stop words)是指在索引和搜索过程中被忽略的常见词汇。这些词汇通常是语言中的功能词,如冠词、介词、连词等,它们在自然语言处理中出现频率很高,但对于文档的语义意义贡献较小。因此,在全文搜索中,停用词通常被过滤掉,以提高搜索性能和相关性。

JavaScript
POST _analyze
{
  "analyzer": "stop",
  "text": "The quick brown fox jumps over the lazy dog."
}
1
2
3
4
5
6

效果:

JavaScript

{
  "tokens" : [
    {
      "token" : "quick",
      "start_offset" : 4,
      "end_offset" : 9,
      "type" : "word",
      "position" : 1
    },
    {
      "token" : "brown",
      "start_offset" : 10,
      "end_offset" : 15,
      "type" : "word",
      "position" : 2
    },
    {
      "token" : "fox",
      "start_offset" : 16,
      "end_offset" : 19,
      "type" : "word",
      "position" : 3
    },
    {
      "token" : "jumps",
      "start_offset" : 20,
      "end_offset" : 25,
      "type" : "word",
      "position" : 4
    },
    {
      "token" : "over",
      "start_offset" : 26,
      "end_offset" : 30,
      "type" : "word",
      "position" : 5
    },
    {
      "token" : "lazy",
      "start_offset" : 35,
      "end_offset" : 39,
      "type" : "word",
      "position" : 7
    },
    {
      "token" : "dog",
      "start_offset" : 40,
      "end_offset" : 43,
      "type" : "word",
      "position" : 8
    }
  ]
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55

上边内置的分词器都不支持对中文分词,我们用标准分词器测试:

JavaScript
POST /_analyze
{
  "analyzer": "standard",
  "text": "黑马程序员学习java太棒了"
}
1
2
3
4
5

结果如下:

JSON
{
  "tokens" : [
    {
      "token" : "黑",
      "start_offset" : 0,
      "end_offset" : 1,
      "type" : "",
      "position" : 0
    },
    {
      "token" : "马",
      "start_offset" : 1,
      "end_offset" : 2,
      "type" : "",
      "position" : 1
    },
    {
      "token" : "程",
      "start_offset" : 2,
      "end_offset" : 3,
      "type" : "",
      "position" : 2
    },
    {
      "token" : "序",
      "start_offset" : 3,
      "end_offset" : 4,
      "type" : "",
      "position" : 3
    },
    {
      "token" : "员",
      "start_offset" : 4,
      "end_offset" : 5,
      "type" : "",
      "position" : 4
    },
    {
      "token" : "学",
      "start_offset" : 5,
      "end_offset" : 6,
      "type" : "",
      "position" : 5
    },
    {
      "token" : "习",
      "start_offset" : 6,
      "end_offset" : 7,
      "type" : "",
      "position" : 6
    },
    {
      "token" : "java",
      "start_offset" : 7,
      "end_offset" : 11,
      "type" : "",
      "position" : 7
    },
    {
      "token" : "太",
      "start_offset" : 11,
      "end_offset" : 12,
      "type" : "",
      "position" : 8
    },
    {
      "token" : "棒",
      "start_offset" : 12,
      "end_offset" : 13,
      "type" : "",
      "position" : 9
    },
    {
      "token" : "了",
      "start_offset" : 13,
      "end_offset" : 14,
      "type" : "",
      "position" : 10
    }
  ]
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82

可以看到,标准分词器智能1字1词条,无法正确对中文做分词。

接下来我们安装一种可以进行中文分词的分词器:IK分词器。

4.5.2 IK分词器 ​

4.5.2.1 安装 ​

IK 分词器是一种广泛应用于中文文本处理的分词工具,尤其在中文搜索引擎和文本分析领域非常常见。

方案一:在线安装

运行一个命令即可:

Shell
docker exec -it es ./bin/elasticsearch-plugin  install https://github.com/medcl/elasticsearch-analysis-ik/releases/download/v7.12.1/elasticsearch-analysis-ik-7.12.1.zip
1

然后重启es容器:

Shell
docker restart es
1

方案二:离线安装

如果网速较差,也可以选择离线安装。

首先,找到Elasticsearch容器的plugins数据卷目录:/data/soft/es7.17.7/xzb/plugins

我们需要把IK分词器上传至这个目录。

找到课前资料提供的ik分词器插件,

配图48

将ik目录上传至虚拟机的/data/soft/es7.17.7/xzb/plugins这个目录:

最后,重启es容器:

Shell
docker restart es
1

我们在安装ES时已经将此目录拷贝到了虚拟机,所以IK分词已经安装成功。

4.5.2.2 测试 ​

下边测试IK分词器:

IK分词器包含两种模式:

ik_smart:智能模式

  • 特点:尽可能地减少输出的词数,适合用于标题或者短文本的分词。
  • 示例:对于输入“中华人民共和国”,智能模式会输出“中华人民共和国”。

ik_max_word:最细粒度模式

  • 特点:尽可能多地输出词,适合用于正文或者长文本的分词。
  • 示例:对于输入“中华人民共和国”,细粒度模式会输出“中华”、“人民”、“中华人”、“中华人民”、“中华人民共和国”。

我们先用智能模式测试:

JSON
POST /_analyze
{
  "analyzer": "ik_smart",
  "text": "中华人民共和国"
}
1
2
3
4
5

执行结果如下:

JSON
{
  "tokens" : [
    {
      "token" : "中华人民共和国",
      "start_offset" : 0,
      "end_offset" : 7,
      "type" : "CN_WORD",
      "position" : 0
    }
  ]
}
1
2
3
4
5
6
7
8
9
10
11

我们再用细粒度模式测试:

JSON
POST /_analyze
{
  "analyzer": "ik_max_word",
  "text": "中华人民共和国"
}
1
2
3
4
5
6
7

执行结果如下:

JSON

{
  "tokens" : [
    {
      "token" : "中华人民共和国",
      "start_offset" : 0,
      "end_offset" : 7,
      "type" : "CN_WORD",
      "position" : 0
    },
    {
      "token" : "中华人民",
      "start_offset" : 0,
      "end_offset" : 4,
      "type" : "CN_WORD",
      "position" : 1
    },
    {
      "token" : "中华",
      "start_offset" : 0,
      "end_offset" : 2,
      "type" : "CN_WORD",
      "position" : 2
    },
    {
      "token" : "华人",
      "start_offset" : 1,
      "end_offset" : 3,
      "type" : "CN_WORD",
      "position" : 3
    },
    {
      "token" : "人民共和国",
      "start_offset" : 2,
      "end_offset" : 7,
      "type" : "CN_WORD",
      "position" : 4
    },
    {
      "token" : "人民",
      "start_offset" : 2,
      "end_offset" : 4,
      "type" : "CN_WORD",
      "position" : 5
    },
    {
      "token" : "共和国",
      "start_offset" : 4,
      "end_offset" : 7,
      "type" : "CN_WORD",
      "position" : 6
    },
    {
      "token" : "共和",
      "start_offset" : 4,
      "end_offset" : 6,
      "type" : "CN_WORD",
      "position" : 7
    },
    {
      "token" : "国",
      "start_offset" : 6,
      "end_offset" : 7,
      "type" : "CN_CHAR",
      "position" : 8
    }
  ]
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69

4.5.3.拓展词典 ​

随着互联网的发展,“造词运动”也越发的频繁。出现了很多新的词语,在原有的词汇列表中并不存在。比如:“泰裤辣”,“传智播客” 等。

IK分词器无法对这些词汇分词,测试一下:

JSON
POST /_analyze
{
  "analyzer": "ik_max_word",
  "text": "传智播客开设大学,真的泰裤辣!"
}
1
2
3
4
5

结果:

JSON
{
  "tokens" : [
    {
      "token" : "传",
      "start_offset" : 0,
      "end_offset" : 1,
      "type" : "CN_CHAR",
      "position" : 0
    },
    {
      "token" : "智",
      "start_offset" : 1,
      "end_offset" : 2,
      "type" : "CN_CHAR",
      "position" : 1
    },
    {
      "token" : "播",
      "start_offset" : 2,
      "end_offset" : 3,
      "type" : "CN_CHAR",
      "position" : 2
    },
    {
      "token" : "客",
      "start_offset" : 3,
      "end_offset" : 4,
      "type" : "CN_CHAR",
      "position" : 3
    },
    {
      "token" : "开设",
      "start_offset" : 4,
      "end_offset" : 6,
      "type" : "CN_WORD",
      "position" : 4
    },
    {
      "token" : "大学",
      "start_offset" : 6,
      "end_offset" : 8,
      "type" : "CN_WORD",
      "position" : 5
    },
    {
      "token" : "真的",
      "start_offset" : 9,
      "end_offset" : 11,
      "type" : "CN_WORD",
      "position" : 6
    },
    {
      "token" : "泰",
      "start_offset" : 11,
      "end_offset" : 12,
      "type" : "CN_CHAR",
      "position" : 7
    },
    {
      "token" : "裤",
      "start_offset" : 12,
      "end_offset" : 13,
      "type" : "CN_CHAR",
      "position" : 8
    },
    {
      "token" : "辣",
      "start_offset" : 13,
      "end_offset" : 14,
      "type" : "CN_CHAR",
      "position" : 9
    }
  ]
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75

可以看到,传智播客和泰裤辣都无法正确分词。

所以要想正确分词,IK分词器的词库也需要不断的更新,IK分词器提供了扩展词汇的功能。

1)打开IK分词器config目录:

配图49

注意,如果采用在线安装的通过,默认是没有config目录的,需要把课前资料提供的ik下的config上传至对应目录。

2)在IKAnalyzer.cfg.xml配置文件内容添加:

JavaScript



        IK Analyzer 扩展配置
        
        
         
        
        
         
        
        words_location -->
1
2
3
4
5
6
7
8
9
10
11
12
13
14

ext_dict:本地扩展词典

ext_stopwords:本地扩展停用词

remote_ext_dict:远程扩展词典,配置一个http连接,通过经连接可以获取扩展词典。

remote_ext_stopwords:远程扩展停用词

我们用ext_dict进行测试,配置一个本地扩展词典文件。

3)在IK分词器的config目录新建一个 ext.dic,一行占一个词语。

Plain
传智播客
泰裤辣
1
2

在IKAnalyzer.cfg.xml文件中配置ext.dic。

配图50

4)重启elasticsearch

Shell
docker restart elasticsearch7.17.7
1

再次测试,请求:

JavaScript
POST /_analyze
{
  "analyzer": "ik_max_word",
  "text": "传智播客开设大学,真的泰裤辣!"
}
1
2
3
4
5

可以发现传智播客和泰裤辣都正确分词了:

JSON
{
  "tokens" : [
    {
      "token" : "传智播客",
      "start_offset" : 0,
      "end_offset" : 4,
      "type" : "CN_WORD",
      "position" : 0
    },
    {
      "token" : "开设",
      "start_offset" : 4,
      "end_offset" : 6,
      "type" : "CN_WORD",
      "position" : 1
    },
    {
      "token" : "大学",
      "start_offset" : 6,
      "end_offset" : 8,
      "type" : "CN_WORD",
      "position" : 2
    },
    {
      "token" : "真的",
      "start_offset" : 9,
      "end_offset" : 11,
      "type" : "CN_WORD",
      "position" : 3
    },
    {
      "token" : "泰裤辣",
      "start_offset" : 11,
      "end_offset" : 14,
      "type" : "CN_WORD",
      "position" : 4
    }
  ]
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40

4.4.4总结 ​

分词器的作用是什么?

  • 创建倒排索引时,对文档分词
  • 用户搜索时,对输入的内容分词

IK分词器有几种模式?

  • ik_smart:智能模式,粗粒度
  • ik_max_word:细粒度模式,细粒度

IK分词器如何拓展词条?

  • 利用config目录的IkAnalyzer.cfg.xml文件添加拓展词典和停用词典
  • 在词典中添加拓展词条或者停用词条

作业 ​

自动取消超时未支付订单 ​

参考 2.3描述完成“自动取消超时未支付订单”功能。

← day04-MQ基础day05-RabbitMQ集群部署 →








如果发现文档内容有错误或排版错乱,请及时联系站长老苗修改,不胜感激。联系方式
关于我们 | 隐私政策 | 豫ICP备2026003386号-4 | 豫公网安备41010202004008号
目录

本页无章节