Masstransit(三)延迟消息
延迟就是发送者延迟设置的时间后,在给消费者发送消息。继上篇的思想,还是需要两个消费者(这个不是一定需要两个消费者,而是我想通过一个消费者转发给另一个消费者实现延迟消息)
配置 需要有下面两个方法
rb.UseDelayedMessageScheduler(); 这里会有坑,下面会写
cfg.AddDelayedMessageScheduler();
service.AddMassTransit(cfg =>
{
cfg.AddDelayedMessageScheduler();
cfg.AddConsumer<OrderConsumer>();
cfg.AddConsumer<FanoutConsumer>();
cfg.AddConsumer<DelayConsumer>();
cfg.UsingRabbitMq((ctx, rb) =>
{
// 配置 RabbitMQ 连接
rb.Host("localhost", h =>
{
h.Username("pony");
h.Password("123456");
});
rb.UseDelayedMessageScheduler();
ListenPriority(rb, ctx);
ListenDirect(rb, ctx);
DelayDirect(rb, ctx);
});
消费者一:原来是直接send给消费者二,现在需要使用ScheduleSend方法实现延迟发送。在最后我会写上遇到的坑,必看。
using MassTransit;
using MassTransit.Scheduling;
using ZR.Common.Model;
using ZR.Model.System;
public class OrderConsumer : IConsumer<SubmitOrder>
{
private readonly ISysUserService _sysUser;
private readonly IBus _bus;
private readonly IMessageScheduler _MessageScheduler;
public OrderConsumer(ISysUserService sysUser, IBus bus, IMessageScheduler messageScheduler)
{
_sysUser = sysUser;
_bus = bus;
_MessageScheduler = messageScheduler;
}
public async Task Consume(ConsumeContext<SubmitOrder> context)
{
var user = new SysUser
{
UserName = context.Message.username,
NickName = "direct",
Password = "123456",
};
if (!string.IsNullOrEmpty(user.UserName))
{
//_sysUser.Add(user);
Console.WriteLine(context.Message.username + "准备发送delay队列");
// var uri = new Uri("rabbitmq://localhost/delay-orders");
// var endpoint = await _bus.GetSendEndpoint(uri);
await _MessageScheduler.ScheduleSend(new Uri("rabbitmq://localhost/delay-orders"),
TimeSpan.FromSeconds(10),
new SubmitOrder { username = "LiYan"});
Console.WriteLine("延迟消息已提交");
// if (endpoint != null)
// Console.WriteLine("发送了消息");
// else Console.WriteLine("未监听到direct");
}
else
return;
}
}
消费者二:接收消息后删除100个用户
using MassTransit;
using ZR.Common.Model;
public class DelayConsumer : IConsumer<SubmitOrder>
{
private readonly ISysUserService _sysUser;
public DelayConsumer(ISysUserService sysUser)
{
_sysUser = sysUser;
}
public async Task Consume(ConsumeContext<SubmitOrder> context)
{
Console.WriteLine($"[DelayConsumer] 收到:{context.Message.username}");
for(int i = 0; i < 100; i++)
{
var user = _sysUser.GetFirst(x => x.NickName == "direct");
if (user != null)
{
await _sysUser.DeleteAsync(user);
Console.WriteLine($"删除了{i}个user");
}
else
Console.WriteLine("删除失败!");
}
}
}
成功

踩坑
1.rb.UseDelayedMessageScheduler();这个方法意思是rabbitmq的延迟插件,但是我的rabbitmq可能没有,否则不能实现延迟发送消息,需要去官网下载
文件名:
rabbitmq_delayed_message_exchange-4.1.0.ez
官网
https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases,然后粘到
D:\RabbitMQ Server\rabbitmq_server-4.1.4\plugins下
2.send和schedulesend的地址不一样,这个要记住。下面是ai的结果,你可以自己研究。
Send 和 ScheduleSend 在 MassTransit 里底层走的是 两条完全不同的路由管线。
| Send | 队列 | rabbitmq://localhost/delay-orders |
| ScheduleSend | 延迟交换机 | exchange:submitorder.delay?type=direct 或 rabbitmq://localhost/.delay |
send:
var uri = new Uri("exchange:submitorder.direct?type=direct");
var endpoint = await _bus.GetSendEndpoint(uri);
await endpoint.Send<SubmitOrder>(new SubmitOrder
{
OrderId = Guid.NewGuid(),
username = "Lisa",
Amount = 500
}, ctx => ctx.SetRoutingKey("addUser"));
、、、、、、、、、、、、、、、、、、、、、、、、
var ep = await bus.GetSendEndpoint(new Uri("rabbitmq://localhost/delay-orders"));
await ep.Send(msg);
-
立刻写进
delay-orders队列; -
无延迟;
-
地址必须是 队列名。
ScheduleSend 的地址 —— 延迟交换机
await scheduler.ScheduleSend(
new Uri("rabbitmq://localhost/delay-orders"), // 最终队列
TimeSpan.FromSeconds(10),
msg);
-
MassTransit 会:
-
把消息发到
.delay交换机(类型x-delayed-message); -
在头里写入
x-delay=10000; -
插件到点后再路由到 真正队列;
-
-
你给的 URI 是最终队列,MassTransit 会自动构造
.delay交换机名字; -
不要手动写成
exchange:...,否则插件找不到队列。
为什么不能用同一个地址?
-
Send 不走插件,没有
x-delay头,RabbitMQ 会立即投递; -
ScheduleSend 必须让 延迟插件先截获,再二次路由到队列;
-
如果地址指错,消息就进错地方,消费者永远收不到。
更多推荐


所有评论(0)