【发布时间】:2021-01-08 14:05:09
【问题描述】:
我正在尝试将 activeMQ 与 NMS (C#) 消费者一起使用来获取消息,进行一些处理,然后通过 HttpClient.PostAsync() 将内容发送到 webserivce,所有这些都在 Windows 服务中运行(通过 Topshelf)。
我正在与之通信的下游系统非常敏感,我正在使用单独的确认,以便我可以通过确认或触发自定义重试(即不是 session.recover)来检查响应并采取相应的行动。
由于下游系统不可靠,我一直在尝试几种不同的方法来降低消费者的吞吐量。我认为我可以通过转换为同步并使用预取来完成此操作,但它似乎没有奏效。
我的理解是,对于异步消费者,预取“限制”永远不会被击中,但使用同步方法,预取队列只会在消息被确认时被吃掉,这意味着我可以调整我的侦听器以以一定速率传递消息下游组件可以处理。
使用一个加载了 100 条消息的队列,并使用侦听器(即异步)启动我的代码,然后我可以成功记录 100 条消息已通过。 当我将其更改为使用 consumer.Receive()(或 ReceiveNoWait)时,我永远不会收到消息。
这是我正在为同步消费者尝试的一个 sn-p,其中包含但已注释掉的 async 选项:
public Worker(LogWriter logger, ServiceConfiguration config, IConnectionFactory connectionFactory, IEndpointClient endpointClient)
{
log = logger;
configuration = config;
this.endpointClient = endpointClient;
connection = connectionFactory.CreateConnection();
connection.RedeliveryPolicy = GetRedeliveryPolicy();
connection.ExceptionListener += new ExceptionListener(OnException);
session = connection.CreateSession(AcknowledgementMode.IndividualAcknowledge);
queue = session.GetQueue(configuration.JmsConfig.SourceQueueName);
consumer = session.CreateConsumer(queue);
// Asynchronous
//consumer.Listener += new MessageListener(OnMessage);
// Synchronous
var message = consumer.Receive(TimeSpan.FromSeconds(5));
while (true)
{
if (!Equals(message, null))
{
OnMessage(message);
}
}
}
public void OnMessage(IMessage message)
{
log.DebugFormat("Message {count} Received. Attempt:{attempt}", message.Properties.GetInt("count"), message.Properties.GetInt("NMSXDeliveryCount"));
message.Acknowledge();
}
【问题讨论】:
-
在获得结束标记之前无法处理 Http。如果您使用的是 HTTP 1.0 流模式,则所有内容都集中在一个块中。 HTTP 1.1 是块模式,其中消息被分解成块。现在在 c# 中,同步模式将阻塞,直到收到所有块。如果您在 1.0 中,异步模式将给出一个响应,但在 1.1 中将给出每个块。
-
@jdweng,您的评论与NMS消费者连接ActiveMQ消息代理接收消息有什么关系?