• Golang RabbitMQ实现的延时队列



    前言

    之前做秒杀商城项目的时候使用到了延时队列来解决订单超时问题,本博客就总结一下Golang是如何利用RabbitMQ实现的延时队列的。

    一、延时队列与应用场景

    延迟队列是一种特殊类型的消息队列,用于在一定时间后将消息投递给消费者。它可以用于处理需要延迟执行的任务或者具有定时特性的业务场景。使用延迟队列可以灵活地控制消息的发送和处理时间,适用于很多场景,如订单超时处理、提醒任务等

    具体应用场景有如下:

    1. 订单取消:当订单生成时,将订单消息发送到延迟队列中,并设置延迟时间为十分钟。消费者在十分钟后接收到订单消息并进行关闭操作。

    2. 店铺商品提醒:在店铺创建时,将提醒消息发送到延迟队列中,并设置延迟时间为十天。消费者在十天后接收到消息并发送提醒通知。

    3. 用户登录提醒:用户注册成功后,将提醒消息发送到延迟队列中,并设置延迟时间为三天。消费者在三天后接收到消息并发送短信提醒。

    4. 退款通知:当用户发起退款时,将通知消息发送到延迟队列中,并设置延迟时间为三天。消费者在三天后接收到消息并通知相关运营人员。

    5. 会议提醒:在会议预定时,将提醒消息发送到延迟队列中,并设置延迟时间为预定时间前十分钟。消费者在指定时间点前十分钟接收到消息并发送会议参加通知。

    通过使用延迟队列,可以在指定的时间点触发任务,避免了轮询的低效方式,并且能够满足大量数据和时效性的需求。这种方法提供了更高的性能和实时性,并有效减轻了系统的负载。
    下图是订单超时处理的流程图。
    在这里插入图片描述

    二、RabbitMQ如何实现延时队列

    虽然 rabbitmq 没有延时队列的功能,但是稍微变动一下也是可以实现的。
    通过设置消息的 TTL 和 DLX 等参数,可以将消息转发到一个指定的队列中,以便在一定的时间后再进行处理。

    实现延时队列的基本要素

    1、存在一个倒计时机制:Time To Live(TTL)
    2、当到达时间点的时候会触发一个发送消息的事件:Dead Letter Exchanges(DLX)

    基于第一点,我利用的是消息存在过期时间这一特性, 消息一旦过期就会变成dead letter,可以让单独的消息过期,也可以设置整个队列消息的过期时间 而rabbitmq会有限取两个值的最小
    **基于第二点,**是用到了rabbitmq的过期消息处理机制: . x-dead-letter-exchange 将过期的消息发送到指定的 exchange 中 . x-dead-letter-routing-key 将过期的消息发送到自定的 route当中

    整体的实现原理如下

    发送者将消息发送到延时队列上并设置过期时间,当过期时间到达时,消息会被自动转发到指定的交换机和队列中供接收者消费。
    1、建立与 RabbitMQ 服务器的连接并创建通道。
    2、发送者通过 ch.Publish 方法将消息发送到延时队列(“test_delay”)上,设置消息的过期时间。
    3、延时队列中的消息在到达过期时间后会自动被发送到 “logs” 交换机,由交换机将消息广播给所有绑定的队列。
    4、接收者通过监听 “test_logs” 队列接收并处理消息。当有消息到达时,会触发回调函数进行处理。
    也就是说要实现延时队列,消费者必须试实现两个队列。
    一个是延时队列(“test_delay”),另一个是接收延时消息的队列(“test_logs”)。

    这两个队列的作用如下:
    延时队列(“test_delay”):这个队列用于接收需要延时发送的消息。发送者通过将消息发送到延时队列,设置消息的过期时间。当消息过期时,RabbitMQ 会自动将消息转发到指定的交换机和队列中。
    接收延时消息的队列(“test_logs”):这个队列用于接收延时消息。在示例中,这个队列是通过将 “test_logs” 队列绑定到 “logs” 交换机上来实现的。交换机会将消息广播给所有绑定的队列,因此当延时消息到达过期时间后,会被发送到这个队列中供消费者进行处理。
    通过使用两个队列,消息可以被延时发送到指定的队列,并在过期后自动转发到接收队列,实现了延时发送和消费的功能。

    三、Go语言实战

    生产者

    首先建立与 RabbitMQ 服务器的连接,并创建一个通道。然后,通过 ch.Publish 方法将消息发送到延时队列上。这里使用的是空字符串作为交换机(exchange),表示不选择任何交换机,只将消息发送到指定的队列(“test_delay”)。在消息的属性中,设置了消息的过期时间为 5 秒。

    
    func main() {
    	conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
    	failOnError(err, "Failed to connect to RabbitMQ")
    	defer conn.Close()
    	ch, err := conn.Channel()
    	failOnError(err, "Failed to open a channel")
    	defer ch.Close()
    	body := "hello"
    	// 将消息发送到延时队列上
    	err = ch.Publish(
    		"", 				// exchange 这里为空则不选择 exchange
    		"test_delay",     	// routing key
    		false,  			// mandatory
    		false,  			// immediate
    		amqp.Publishing{
    			ContentType: "text/plain",
    			Body:        []byte(body),
    			Expiration: "5000",	// 设置五秒的过期时间
    		})
    	failOnError(err, "Failed to publish a message")
    
    	log.Printf(" [x] Sent %s", body)
    }
    
    
    • 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

    消费者

    同样建立与 RabbitMQ 服务器的连接,并创建一个通道。然后,声明了一个名为 “logs” 的交换机,类型为 “fanout”,并且可持久化,表示该交换机会将消息广播给所有绑定的队列。接着,声明了一个常规的队列 “test_logs”,并将其绑定到 “logs” 交换机上。之后,声明了一个延时队列 “test_delay”,并设置了该队列的 x-dead-letter-exchange 参数为 “logs”,即当消息过期时将消息发送到 “logs” 交换机。最后,通过 ch.Consume 方法监听 “test_logs” 队列,并在回调函数中处理接收到的消息。

    
    func main() {
    	// 建立链接
    	conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
    	failOnError(err, "Failed to connect to RabbitMQ")
    	defer conn.Close()
    
    	ch, err := conn.Channel()
    	failOnError(err, "Failed to open a channel")
    	defer ch.Close()
    
    	// 声明一个主要使用的 exchange
    	err = ch.ExchangeDeclare(
    		"logs",   // name
    		"fanout", // type
    		true,     // durable
    		false,    // auto-deleted
    		false,    // internal
    		false,    // no-wait
    		nil,      // arguments
    	)
    	failOnError(err, "Failed to declare an exchange")
    
    	// 声明一个常规的队列, 其实这个也没必要声明,因为 exchange 会默认绑定一个队列
    	q, err := ch.QueueDeclare(
    		"test_logs",    // name
    		false, // durable
    		false, // delete when unused
    		true,  // exclusive
    		false, // no-wait
    		nil,   // arguments
    	)
    	failOnError(err, "Failed to declare a queue")
    
        /**
         * 注意,这里是重点!!!!!
         * 声明一个延时队列, ß我们的延时消息就是要发送到这里
         */
    	_, errDelay := ch.QueueDeclare(
    		"test_delay",    // name
    		false, // durable
    		false, // delete when unused
    		true,  // exclusive
    		false, // no-wait
    		amqp.Table{
    			// 当消息过期时把消息发送到 logs 这个 exchange
    			"x-dead-letter-exchange":"logs",
    		},   // arguments
    	)
    	failOnError(errDelay, "Failed to declare a delay_queue")
    
    	err = ch.QueueBind(
    		q.Name, // queue name, 这里指的是 test_logs
    		"",     // routing key
    		"logs", // exchange
    		false,
    		nil)
    	failOnError(err, "Failed to bind a queue")
    
    	// 这里监听的是 test_logs
    	msgs, err := ch.Consume(
    		q.Name, // queue name, 这里指的是 test_logs
    		"",     // consumer
    		true,   // auto-ack
    		false,  // exclusive
    		false,  // no-local
    		false,  // no-wait
    		nil,    // args
    	)
    	failOnError(err, "Failed to register a consumer")
    
    	forever := make(chan bool)
    
    	go func() {
    		for d := range msgs {
    			log.Printf(" [x] %s", d.Body)
    		}
    	}()
    
    	log.Printf(" [*] Waiting for logs. To exit press CTRL+C")
    	<-forever
    }
    
    • 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
  • 相关阅读:
    【编程不良人】Redis后端实战学习笔记01---NoSQL简介、分类、应用场景;Redis特点、安装、相关指令
    升降 ubuntu 内核的好工具 mainline
    内核同步机制-自旋锁(spin_lock)
    Apipost-Helper:IDEA中的类postman工具
    【ArcGIS Pro微课1000例】0032:创建具有指定高程Z值的矢量数据
    刷爆力扣之子数组最大平均数 I
    基于BDD的接口自动化框架开箱即用
    KSO - .net6项目中使用RabbitMQ实际项目代码和思路讲解,包括各种踩坑
    JS - 防抖与节流
    软件设计师——信息安全知识
  • 原文地址:https://blog.csdn.net/weixin_44545838/article/details/132655899