延迟就是发送者延迟设置的时间后,在给消费者发送消息。继上篇的思想,还是需要两个消费者(这个不是一定需要两个消费者,而是我想通过一个消费者转发给另一个消费者实现延迟消息)

配置 需要有下面两个方法

 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的结果,你可以自己研究。

SendScheduleSend 在 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 会:

    1. 把消息发到 .delay 交换机(类型 x-delayed-message);

    2. 在头里写入 x-delay=10000

    3. 插件到点后再路由到 真正队列

  • 你给的 URI 是最终队列,MassTransit 会自动构造 .delay 交换机名字;

  • 不要手动写成 exchange:...,否则插件找不到队列。

 为什么不能用同一个地址?

  • Send 不走插件,没有 x-delay,RabbitMQ 会立即投递;

  • ScheduleSend 必须让 延迟插件先截获,再二次路由到队列;

  • 如果地址指错,消息就进错地方,消费者永远收不到。

Logo

开源鸿蒙跨平台开发社区汇聚开发者与厂商,共建“一次开发,多端部署”的开源生态,致力于降低跨端开发门槛,推动万物智联创新。

更多推荐