You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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,提问作者Николай Дойчев

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 08:40:48