您的位置:首页 > 编程语言 > C#

C#通过rabbitmq实现定时任务(延时队列)

2021-04-26 17:27 1436 查看


本文主要讲解如何通过RabbitMQ实现定时任务(延时队列)

环境准备

需要在MQ中进行安装插件 地址链接
插件介绍地址:https://www.rabbitmq.com/blog/2015/04/16/scheduling-messages-with-rabbitmq/

使用场景

作为一个新的预支付订单被初始化放置,如果该订单在指定时间内未进行支付,则将被认为超时订单进行关闭处理;电商系统中应用较多,用户购买商品产生订单,但未进行支付,订单产生30分钟内未支付将关闭订单(且满足该场景数量庞大),不可能采用人工干预。

代码介绍

生产者

var factory = new ConnectionFactory()
{
Uri = new Uri("MQ地址")
};

using var connection = factory.CreateConnection();
using var channel = connection.CreateModel();

var exchangeName = "delay-exchange";
var routingkey = "delay.delay";
var queueName = "delay_queueName";
//设置Exchange队列类型
var argMaps = new Dictionary<string, object>()
{
{"x-delayed-type", "topic"}
};
//设置当前消息为延时队列
channel.ExchangeDeclare(exchange: exchangeName, type: "x-delayed-message", true, false, argMaps);
channel.QueueDeclare(queueName, true, false, false, argMaps);
channel.QueueBind(queueName, exchangeName, routingkey);
for (int i = 0; i < 3; i++)
{
var time = 1000 * 5;
var message = $@"发送时间为 {DateTime.Now:yyyy-MM-dd HH:mm:ss} 延时时间为:{time}";
var body = Encoding.UTF8.GetBytes(message);
var props = channel.CreateBasicProperties();
//设置消息的过期时间
props.Headers = new Dictionary<string, object>()
{
{  "x-delay", 5000 }
};
channel.BasicPublish(exchange: exchangeName,
routingKey: routingkey,
basicProperties: props,
body: body);
Console.WriteLine(message);

}
Console.ReadLine();

消费者(自动绑定队列写法)

var factory = new ConnectionFactory()
{
Uri = new Uri(MQ地址)
};
using var connection = factory.CreateConnection();
using var channel = connection.Cr
56c
eateModel();
var queueName = "delay_queueName";
channel.QueueDeclare(queueName, true, false, false, null);
var consumer = new EventingBasicConsumer(channel);
consumer.Received += (model, ea) =>
{
var body = ea.Body;
var message = Encoding.UTF8.GetString(body);
var routingKey = ea.RoutingKey;
Console.WriteLine($@"接受到消息的时间为 {DateTime.Now:yyyy-MM-dd HH:mm:ss},routingKey:{routingKey} message:{message} ");
};
channel.BasicConsume(queue: queueName,
autoAck: true,
consumer: consumer);
Console.ReadLine();

消费者(手动绑定队列写法)

var factory = new ConnectionFactory()
{
Uri = new Uri(MQ地址)
};
using var connection = factory.CreateConnection();
using var channel = connection.CreateModel();
var exchangeName = "delay-exchange";
var routingkey = "delay.delay";
var queueName = "delay_queueName";
var autoDelete = true;

var argMaps = new Dictionary<string, object>()
{
{"x-delayed-type", "topic"}
};
channel.ExchangeDeclare(exchange: exchangeName
ad8
, type: "x-delayed-message", true, false, argMaps);
channel.QueueDeclare(queueName, true, false, false, argMaps);
channel.QueueBind(queue: queueName, exchange: exchangeName, routingKey: routingkey);
//channel.QueueDeclare(queueName, true, false, false, null);
var consumer = new EventingBasicConsumer(channel);
consumer.Received += (model, ea) =>
{
var body = ea.Body;
var message = Encoding.UTF8.GetString(body);
var routingKey = ea.RoutingKey;
Console.WriteLine($@"接受到消息的时间为 {DateTime.Now:yyyy-MM-dd HH:mm:ss},routingKey:{routingKey} message:{message} ");
};
channel.BasicConsume(queue: queueName,
autoAck: true,
consumer: consumer);
Console.ReadLine();

最终实现效果(两个消费者)

在上述实现中,其实主要靠以下参数来帮我们实现当前功能

声明Exchange中的 type: "x-delayed-message" 这个表明当前队列为延时消息队列
声明Exchange中arguments中的 {"x-delayed-type", "topic"} 当前表明当前队列为Topic模式
最后 我们在CreateBasicProperties的Header中设置 { "x-delay", 5000 }来达到消息延时的功能(单位为ms)

建议

如果使用当前模式来做定时任务,在要求消息不丢失的前提下,需要运维同学提供稳定的MQ环境

如有哪里讲得不是很明白或是有错误,欢迎指正
如您喜欢的话不妨点个赞收藏一下吧🙂

内容来自用户分享和网络整理,不保证内容的准确性,如有侵权内容,可联系管理员处理 点击这里给我发消息
标签: