【问题标题】:Scaling Out Subscribers with MassTransit Request/Response conversation使用 MassTransit 请求/响应对话扩展订阅者
【发布时间】:2017-04-26 06:34:22
【问题描述】:

我有一个使用 Masstransit / RabbitMq 实现的企业服务总线,用于 Web 项目。我使用 Masstransit 的 Request/Response 模式在 MQ 上执行 RPC。

我在处理不同类型消息的总线中创建多个 ReceiveEndpoint。 ESB 也使用自定义配置文件来创建多个总线,因此我可以以某种方式在逻辑上实现 Qos,甚至可以通过网络在单独的服务器上分发消费者,以或多或少地提高性能。

最近我看到我的一位消费者被屏蔽了。如果我将一些消息发送给一个长时间运行的消费者(每个需要大约 5 秒的工作时间),看起来就像一个线程响应了我的请求,并发发送的消息等待之前的消息被使用。

到目前为止,我设置 UseConcurrencyLimit(64) 没有任何改变,尝试将 PrefetchCount 增加到 50,但是 RabbitMq 在队列详细信息中将其显示为 0。

为什么消费者一次处理一条消息?

MassTransit v3.5.7

编辑:我后来找到了this thread。对我来说,这似乎是同样的问题。

更新 与 Chris 的示例不同的是,我使用 RabbitMq 并使用反射来创建使用者,这样我就可以使用配置文件管理 ESB。仍然在 RabbitMq 管理控制台上看到 Prefetch Count 0。

var busControl = Bus.Factory.CreateUsingRabbitMq(x =>
        {
            x.UseConcurrencyLimit(64);
            IRabbitMqHost host = x.Host(new Uri(Config.host), h =>
            {
                h.Username("userName");
                h.Password("password");
            });

            var obj = Activator.CreateInstance([SomeExistingType]);

            x.ReceiveEndpoint(host, "queueName", e =>
            {
                e.PrefetchCount = 64;//config.prefetchCount;
                e.PurgeOnStartup = true;
                if (config.retryCount > 0)
                {
                    e.UseRetry(retry => retry.Interval(config.retryCount, config.retryInterval));
                }
                e.SetQueueArgument("x-expires", config.timeout * 1000 /*Seconds to milliseconds*/);
                e.SetQueueArgument("temporary", config.temporary);

                e.Consumer(consumer, f => obj);

            }); 

        })

busControl.StartAsync();

更新 2 当我将预取计数设置为 1 以让 MassTransit 处理工作负载时,因为我有多个消费者服务器和一个 RabbitMq 集群。但是,当我将许多消息发送到使所有线程保持忙碌的队列时,新请求在发送者队列中等待由空闲线程拾取。我增加了预取计数,这样我最终得到了更多用于新请求的空闲线程。按照 Chris 的建议,我在配置接收端点时设置了预取计数。谢谢。 Chris 确认后,我会将 Chris 的回复标记为答复。

【问题讨论】:

  • 在您的示例中,您将单个消费者用于所有消息。通常不是最好的选择,但如果你坚持下去,请确保你的消费者没有任何实例变量,因为它们将被所有并发消息共享。
  • 您可以通过将 Activator.CreateInstance 调用移动到 e.Consumer() 调用的 lambda 方法来解决此问题。
  • @ChrisPatterson 非常感谢,我之前从您的回复中读到了该评论家。在示例中,我觉得我需要修改方法以使这个特定问题更加简单。实际上,我为每种消息类型创建了许多 ReceiveEnpoint。此外,我创建了几条总线来将一些消息组合在一起,因为其中一些使用请求/响应模式,而另一些用于发布/订阅。

标签: rabbitmq masstransit


【解决方案1】:

如果您在 RabbitMQ 控制台中看到 0 以获取预取计数,则您可能将其配置在错误的位置。你应该有一些类似的东西:

cfg.ReceiveEndpoint("my-queue", x =>
{
    x.PrefetchCount = 64;
    x.Consumer<MyConsumer>();
});

这会将接收端点消费者配置为 64 的预取计数,这应该会显示在 RabbitMQ 管理控制台的消费者中。

另一方面,如果您的消费者使用Thread.Sleep() 之类的东西阻塞线程,我已经看到了可能限制线程池并发性的情况。

【讨论】:

  • 感谢您的回答。我已经更新了我的原始问题,请您检查我是否遗漏了什么?
  • 另一个更新:尽管 RabbitMq 显示 Prefetch 计数为 0,但实际上不是 0。我可以说,经过一些测试。你的回答看起来对我有用,我会解释我在原来的问题中看到的错误。请确认...
猜你喜欢
  • 2016-09-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-05-15
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多