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

部署在AWS的ActiveMQ连接数秒后断开,需实现持续消息监听

Fixing Apache ActiveMQ Consumer to Stay Connected and Listen Continuously

Got it, let's get your consumer working properly so it stays connected and processes all incoming messages instead of exiting right away. Here's a breakdown of the issues in your current code and the updated version to fix them:

What's Wrong With Your Current Code?

  • Immediate Exit After Startup: After calling connection.Start(), your code immediately checks if message is null and exits. Most of the time, the message hasn't been received yet (since the listener runs asynchronously), so it prints "No message received!" and terminates, dropping the connection.
  • Single Message Limitation: The static message variable only holds one message—you can't process multiple incoming messages with this setup.
  • No Lifecycle Management: There's no mechanism to keep the main thread alive, so the program exits as soon as the initial check completes.

Updated Code for Continuous Listening

using Apache.NMS;
using Apache.NMS.Util;
using System;

namespace ApacheMQAsync
{
    class Program
    {
        public static void Main(string[] args)
        {
            Uri connecturi = new Uri("URL:61617");
            Console.WriteLine("About to connect to " + connecturi);

            // Use using statements to ensure proper resource cleanup
            using (IConnectionFactory factory = new Apache.NMS.ActiveMQ.ConnectionFactory(connecturi))
            using (IConnection connection = factory.CreateConnection("username", "password"))
            using (ISession session = connection.CreateSession())
            {
                IDestination destination = SessionUtil.GetDestination(session, "queue://FOO.BAR");
                Console.WriteLine("Using destination: " + destination);

                // Create a consumer and attach the listener
                using (IMessageConsumer consumer = session.CreateConsumer(destination))
                {
                    consumer.Listener += OnMessage;
                    connection.Start();

                    Console.WriteLine("Consumer is running. Press any key to exit...");
                    // Block the main thread to keep the program alive
                    Console.ReadLine();
                }
            }
        }

        protected static void OnMessage(IMessage receivedMsg)
        {
            if (receivedMsg is ITextMessage textMessage)
            {
                Console.WriteLine("Received message with ID: " + textMessage.NMSMessageId);
                Console.WriteLine("Received message with text: " + textMessage.Text);
                textMessage.Acknowledge();
            }
            else
            {
                Console.WriteLine($"Received non-text message of type: {receivedMsg.GetType().Name}");
                receivedMsg.Acknowledge();
            }
        }
    }
}

Key Changes Explained

  • Removed Static message Variable: Instead of storing messages in a static field, we process each message directly in the OnMessage callback. This lets us handle an unlimited number of incoming messages.
  • Added Main Thread Blocking: Console.ReadLine() keeps the main thread running, so the program doesn't exit immediately. The connection stays active as long as the program is running.
  • Using using Statements: These ensure that connections, sessions, and consumers are properly disposed when the program exits (when you press a key), preventing resource leaks in ActiveMQ.
  • Improved Message Handling: We added a check to handle non-text messages gracefully, making the consumer more robust.

Now your consumer will stay connected, listen for all incoming messages, and process each one as it arrives. Just press any key in the console when you want to stop it.

内容的提问来源于stack exchange,提问作者Mubashir Koul

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:54:39