首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >MassTransit / RabbitMQ -为什么会跳过这么多消息?

MassTransit / RabbitMQ -为什么会跳过这么多消息?
EN

Stack Overflow用户
提问于 2019-05-08 18:14:54
回答 3查看 3.1K关注 0票数 1

我正在使用.NET /RabbitMQ的生产者/消费者场景中的2个核心控制台应用程序。我需要确保,即使没有消费者启动和运行,来自生产者的消息仍然成功排队。这似乎不适用于Publish() --消息只是消失了,所以我使用Send()代替。消息至少会排队,但是如果没有任何用户运行这些消息,所有消息都会在"_skipped“队列中结束。

这就是我的第一个问题:这是基于需求的正确方法(即使没有消费者正在运行,来自生产者的消息仍在成功排队)?

使用Send(),我的使用者确实可以工作,但是仍然有许多消息从漏洞中掉下来,被抛到"_skipped“队列中。消费者的逻辑是最小的(目前只是记录消息),所以它不是一个长期运行的过程。

这就是我的第二个问题:为什么仍有这么多消息被倾倒到"_skipped“队列中?

这就引出了我的第三个问题:这是否意味着我的消费者也需要收听"_skipped“队列?

我不知道您需要为这个问题看什么代码,但是下面是RabbitMQ管理UI的屏幕截图:

生产者配置:

代码语言:javascript
复制
    static IHostBuilder CreateHostBuilder(string[] args)
    {
        return Host.CreateDefaultBuilder()
                      .ConfigureServices((hostContext, services) =>
                      {
                          services.Configure<ApplicationConfiguration>(hostContext.Configuration.GetSection(nameof(ApplicationConfiguration)));

                          services.AddMassTransit(cfg =>
                          {
                              cfg.AddBus(ConfigureBus);
                          });

                          services.AddHostedService<CardMessageProducer>();
                      })
                      .UseConsoleLifetime()
                      .UseSerilog();
    }

    static IBusControl ConfigureBus(IServiceProvider provider)
    {
        var options = provider.GetRequiredService<IOptions<ApplicationConfiguration>>().Value;

        return Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            var host = cfg.Host(new Uri(options.RabbitMQ_ConnectionString), h =>
            {
                h.Username(options.RabbitMQ_Username);
                h.Password(options.RabbitMQ_Password);
            });

            cfg.ReceiveEndpoint(host, typeof(CardMessage).FullName, e =>
            {
                EndpointConvention.Map<CardMessage>(e.InputAddress);
            });
        });
    }

生产者代码:

代码语言:javascript
复制
Bus.Send(message);

消费者配置:

代码语言:javascript
复制
    static IHostBuilder CreateHostBuilder(string[] args)
    {
        return Host.CreateDefaultBuilder()
                      .ConfigureServices((hostContext, services) =>
                      {
                          services.AddSingleton<CardMessageConsumer>();

                          services.Configure<ApplicationConfiguration>(hostContext.Configuration.GetSection(nameof(ApplicationConfiguration)));

                          services.AddMassTransit(cfg =>
                          {
                              cfg.AddBus(ConfigureBus);
                          });

                          services.AddHostedService<MassTransitHostedService>();
                      })
                      .UseConsoleLifetime()
                      .UseSerilog();
    }

    static IBusControl ConfigureBus(IServiceProvider provider)
    {
        var options = provider.GetRequiredService<IOptions<ApplicationConfiguration>>().Value;

        return Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            var host = cfg.Host(new Uri(options.RabbitMQ_ConnectionString), h =>
            {
                h.Username(options.RabbitMQ_Username);
                h.Password(options.RabbitMQ_Password);
            });

            cfg.ReceiveEndpoint(host, typeof(CardMessage).FullName, e =>
            {
                e.Consumer<CardMessageConsumer>(provider);
            });

            //cfg.ReceiveEndpoint(host, typeof(CardMessage).FullName + "_skipped", e =>
            //{
            //    e.Consumer<CardMessageConsumer>(provider);
            //});
        });
    }

消费者代码:

代码语言:javascript
复制
class CardMessageConsumer : IConsumer<CardMessage>
{
    private readonly ILogger<CardMessageConsumer> logger;
    private readonly ApplicationConfiguration configuration;
    private long counter;

    public CardMessageConsumer(ILogger<CardMessageConsumer> logger, IOptions<ApplicationConfiguration> options)
    {
        this.logger = logger;
        this.configuration = options.Value;
    }

    public async Task Consume(ConsumeContext<CardMessage> context)
    {
        this.counter++;

        this.logger.LogTrace($"Message #{this.counter} consumed: {context.Message}");
    }
}
EN

回答 3

Stack Overflow用户

发布于 2019-05-08 19:00:11

在MassTransit中,_skipped队列是死信队列概念的实现。信息到达那里是因为它们不会被消耗掉。

带有RMQ的MassTransit总是将消息传递给交换,而不是传递到队列。默认情况下,每个MassTransit端点创建一个具有端点名称的队列(如果没有现有队列),创建一个具有相同名称的交换,并将它们绑定到一起。当应用程序具有已配置的使用者(或处理程序)时,还会创建该消息类型的交换(使用消息类型作为交换名称),并将端点交换绑定到消息类型exchange。因此,当您使用Publish时,消息将被发布到消息类型exchange中,并相应地使用端点绑定(或多个绑定)来传递。当您使用Send时,没有使用消息类型交换,因此消息直接到达目标交换。而且,正如@maldworth正确地指出的那样,每个MassTransit端点只期望获得它可以使用的消息。如果它不知道如何使用该消息-消息将被移动到死信队列中。这以及有害消息队列是消息传递的基本模式。

如果您需要排队等待稍后使用,最好的方法是设置连接,但是端点本身(我指的是应用程序)不应该运行。一旦应用程序启动,它将使用所有排队的消息。

票数 2
EN

Stack Overflow用户

发布于 2019-05-08 18:43:05

当使用者启动总线bus.Start()时,它所做的事情之一就是为传输创建所有的交换和队列。如果您需要在使用者之前发布/发送,那么您唯一的选择就是运行DeployTopologyOnly。不幸的是,这个特性没有记录在正式文档中,但是单元测试在这里:Specs.cs

当消息被发送到不知道如何处理的使用者时,跳过的队列就会发生。

例如,如果您有一个可以处理IConsumer<MyMessageA>的使用者,它位于接收端点名称“my a”上。但是您的消息生成程序会执行Send<MyMessageB>(Uri("my-queue-a")...),这是一个问题。消费者只理解A,它不知道如何处理B,所以它只是将它移动到跳过的队列中,然后继续。

票数 0
EN

Stack Overflow用户

发布于 2022-10-16 15:34:36

在我的例子中,同一个队列同时侦听多个消费者。

票数 0
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/56046802

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档