如何在Apache.NMS.ActiveMQ中实现消息并发处理?
.NET7下Apache.NMS.ActiveMQ实现并发消息消费
问题描述
我在.NET 7项目中使用Apache.NMS.ActiveMQ 1.7.2,当前消息消费为串行模式(处理完一条消息后再处理下一条),希望实现多消息并发消费,让每条消息在独立线程中独立处理。以下是我的初始化代码,求现成示例。
原消费代码
public class AmqClientService : IMessageProvider { private readonly ILogger _logger; private readonly IMessageProcessor _messageProcessor; private readonly AmqClientSettings _settings = new AmqClientSettings(); private IConnection connection; //private ISession session; private IMessageConsumer consumer; public AmqClientService(ILogger logger, IConfiguration config, IMessageProcessor messageProcessor) { _logger = logger; _messageProcessor = messageProcessor; config.Bind(nameof(AmqClientSettings), _settings); } public async Task WaitAsync(CancellationToken stoppingToken) { var exampleProvider = "tcp://localhost:61616"; _logger.Information($"initialize AMQ connection with provider {exampleProvider}"); Uri connecturi = new Uri(exampleProvider); ConnectionFactory connectionFactory = new ConnectionFactory(connecturi); this.connection = connectionFactory.CreateConnection(); using var session = connection.CreateSession(); this.connection.Start(); using (var consumer = session.CreateConsumer(SessionUtil.GetDestination(session, AppSettings.QueueName))) { while (!stoppingToken.IsCancellationRequested) { _logger.Information($"Waiting for SAP message"); var message = consumer.Receive() as ITextMessage; _logger.Information($"Received message {message.Text}, for {message.NMSDestination}"); if (await _messageProcessor.TryProcessAsync(message.Text)) { message.Acknowledge(); } } } } }
原DI配置
serviceCollection.AddMemoryCache(); serviceCollection.AddSingleton(Log.Logger); serviceCollection.AddSingleton<IMessageProcessor, MessageProcessor>(); serviceCollection.AddSingleton<IMessageProvider, AmqClientService>(); serviceCollection.AddHostedService<Worker>(); serviceCollection.BuildServiceProvider();
解决方案
核心思路
- 改用异步接收方法
ReceiveAsync()避免阻塞主线程 - 将消息处理逻辑封装到独立任务中实现并发
- 注意
ISession非线程安全的特性,合理规划Session的使用范围
修改后的消费代码(基础并发版)
public class AmqClientService : IMessageProvider { private readonly ILogger _logger; private readonly IMessageProcessor _messageProcessor; private readonly AmqClientSettings _settings = new AmqClientSettings(); // 可选:限制并发数量,根据服务器配置调整 private readonly SemaphoreSlim _concurrencySemaphore = new SemaphoreSlim(5); private IConnection connection; public AmqClientService(ILogger logger, IConfiguration config, IMessageProcessor messageProcessor) { _logger = logger; _messageProcessor = messageProcessor; config.Bind(nameof(AmqClientSettings), _settings); } public async Task WaitAsync(CancellationToken stoppingToken) { var exampleProvider = "tcp://localhost:61616"; _logger.Information($"初始化AMQ连接,地址:{exampleProvider}"); Uri connecturi = new Uri(exampleProvider); ConnectionFactory connectionFactory = new ConnectionFactory(connecturi); // 设置客户端确认模式,允许在独立线程中确认消息 connectionFactory.AcknowledgeMode = AcknowledgeMode.ClientAcknowledge; this.connection = connectionFactory.CreateConnection(); this.connection.Start(); // 接收用Session单线程使用,保证线程安全 using var receiveSession = connection.CreateSession(); using var consumer = receiveSession.CreateConsumer(SessionUtil.GetDestination(receiveSession, AppSettings.QueueName)); _logger.Information("开始监听消息队列"); while (!stoppingToken.IsCancellationRequested) { try { // 异步接收,避免阻塞循环 var message = await consumer.ReceiveAsync(TimeSpan.FromSeconds(1)) as ITextMessage; if (message == null) continue; _logger.Information($"收到消息:{message.Text},目标队列:{message.NMSDestination}"); // 启动独立任务处理消息,传递取消令牌 _ = ProcessMessageAsync(message, stoppingToken); } catch (Exception ex) { _logger.Error(ex, "监听消息队列时出现异常"); await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken); } } _logger.Information("停止监听消息队列"); connection.Close(); } private async Task ProcessMessageAsync(ITextMessage message, CancellationToken stoppingToken) { await _concurrencySemaphore.WaitAsync(stoppingToken); try { if (await _messageProcessor.TryProcessAsync(message.Text)) { message.Acknowledge(); _logger.Information($"消息处理完成,已确认:{message.Text}"); } else { _logger.Warning($"消息处理失败:{message.Text}"); } } catch (Exception ex) { _logger.Error(ex, $"处理消息时出错:{message.Text}"); // 处理失败可根据需求选择是否重新入队或进入死信队列 } finally { _concurrencySemaphore.Release(); } } }
更安全的多Session版本(适合处理逻辑需MQ交互场景)
如果你的消息处理逻辑中需要和MQ进行额外交互(比如发送响应消息),建议为每个处理任务创建独立Session,避免线程安全问题:
private async Task ProcessMessageAsync(ITextMessage message, CancellationToken stoppingToken) { await _concurrencySemaphore.WaitAsync(stoppingToken); ISession processingSession = null; try { // 创建独立Session用于消息处理 processingSession = connection.CreateSession(); if (await _messageProcessor.TryProcessAsync(message.Text)) { message.Acknowledge(); _logger.Information($"消息处理完成,已确认:{message.Text}"); } else { _logger.Warning($"消息处理失败:{message.Text}"); } } catch (Exception ex) { _logger.Error(ex, $"处理消息时出错:{message.Text}"); } finally { processingSession?.Close(); _concurrencySemaphore.Release(); } }
关键注意事项
- Session线程安全:
ISession不能在多线程间共享,接收消息的Session需单线程使用,处理逻辑如需Session必须单独创建 - 并发控制:通过
SemaphoreSlim限制并发数量,避免系统资源耗尽 - 消息确认:仅在处理成功后调用
Acknowledge(),处理失败时消息会自动重新入队(需配置ActiveMQ的重发策略和死信队列) - 取消令牌传递:所有异步操作都需传递
stoppingToken,确保应用停止时能终止所有任务
DI配置保持不变
原DI配置无需修改,继续使用即可:
serviceCollection.AddMemoryCache(); serviceCollection.AddSingleton(Log.Logger); serviceCollection.AddSingleton<IMessageProcessor, MessageProcessor>(); serviceCollection.AddSingleton<IMessageProvider, AmqClientService>(); serviceCollection.AddHostedService<Worker>(); serviceCollection.BuildServiceProvider();
内容的提问来源于stack exchange,提问作者Николай Дойчев
相关产品推荐
相关产品推荐

