【发布时间】: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