部署在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 ifmessageis 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
messagevariable 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
messageVariable: Instead of storing messages in a static field, we process each message directly in theOnMessagecallback. 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
usingStatements: 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
相关产品推荐
相关产品推荐

