【问题标题】:MQQueueManager message poolingMQQueueManager 消息池
【发布时间】:2018-08-02 01:16:52
【问题描述】:

我过去曾将 RabbitMq 用作 MessageQueue,在收到消息时触发事件非常简单。

我查看了 IBM 安装程序提供的 .NET 源代码,但我发现处理它的方法并不好。查看示例 SimpleSubscribe 它做了这样的事情来池

// getting messages continuously
for (int i = 1; i <= numberOfMsgs; i++)
{
    // creating a message object
    message = new MQMessage();

    try
    {
        topic.Get(message);
        Console.WriteLine("Message " + i + " got = " + message.ReadString(message.MessageLength));
        message.ClearMessage();
    }
    catch (MQException mqe)
    {
        if (mqe.ReasonCode == 2033)
        {
            ++time;
            --i;
            Console.WriteLine("No message available");
            Thread.Sleep(1000);
            //waiting for 10sec
            if (time > 10)
            {
                Console.WriteLine("Timeout : No message available");
                break;
            }
            continue;
        }
        else
        {
            Console.WriteLine("MQException caught: {0} - {1}", mqe.ReasonCode, mqe.Message);
        }
    }
}

其中numberOfMsgs 的数量是作为参数传递的整数。如果我通过 -1,它可能会在无限处汇集,但我认为这不是一个好方法。这导致我以特定的时间间隔汇集队列,然后可能会以 100% 的速度旋转一个线程我没有放 Thread.Sleep 有没有人有更好的方法?

【问题讨论】:

    标签: c# .net ibm-mq


    【解决方案1】:

    您应该使用“MQGET with Wait”选项而不是使用轮询。 MQGET 调用将等待消息到达并在收到消息时立即返回,否则将等待“x”毫秒。

    这是一个简单的 C# MQ 程序:

    using System;
    using System.Collections;
    using System.Collections.Generic;
    using System.Text;
    using IBM.WMQ;
    
    /// <summary> Program Name
    /// MQTest62
    ///
    /// Description
    /// This C# class will connect to a remote queue manager
    /// and get messages from a queue using a managed .NET environment.
    ///
    /// Sample Command Line Parameters
    /// -h 127.0.0.1 -p 1415 -c TEST.CHL -m MQWT1 -q TEST.Q1 -u tester -x mypwd
    /// </summary>
    /// <author>  Roger Lacroix
    /// </author>
    namespace MQTest62
    {
       public class MQTest62
       {
          private Hashtable inParms = null;
          private Hashtable qMgrProp = null;
          private System.String qManager;
          private System.String inputQName;
    
          /*
          * The constructor
          */
          public MQTest62()
              : base()
          {
          }
    
          /// <summary> Make sure the required parameters are present.</summary>
          /// <returns> true/false
          /// </returns>
          private bool allParamsPresent()
          {
             bool b = inParms.ContainsKey("-h") && inParms.ContainsKey("-p") &&
                      inParms.ContainsKey("-c") && inParms.ContainsKey("-m") &&
                      inParms.ContainsKey("-q");
             if (b)
             {
                try
                {
                   System.Int32.Parse((System.String)inParms["-p"]);
                }
                catch (System.FormatException e)
                {
                   b = false;
                }
             }
    
             return b;
          }
    
          /// <summary> Extract the command-line parameters and initialize the MQ variables.</summary>
          /// <param name="args">
          /// </param>
          /// <throws>  IllegalArgumentException </throws>
          private void init(System.String[] args)
          {
             inParms = System.Collections.Hashtable.Synchronized(new System.Collections.Hashtable(14));
             if (args.Length > 0 && (args.Length % 2) == 0)
             {
                for (int i = 0; i < args.Length; i += 2)
                {
                   inParms[args[i]] = args[i + 1];
                }
             }
             else
             {
                throw new System.ArgumentException();
             }
    
             if (allParamsPresent())
             {
                qManager = ((System.String)inParms["-m"]);
                inputQName = ((System.String)inParms["-q"]);
    
                qMgrProp = new Hashtable();
                qMgrProp.Add(MQC.TRANSPORT_PROPERTY, MQC.TRANSPORT_MQSERIES_MANAGED);
    
                qMgrProp.Add(MQC.HOST_NAME_PROPERTY, ((System.String)inParms["-h"]));
                qMgrProp.Add(MQC.CHANNEL_PROPERTY, ((System.String)inParms["-c"]));
    
                try
                {
                   qMgrProp.Add(MQC.PORT_PROPERTY, System.Int32.Parse((System.String)inParms["-p"]));
                }
                catch (System.FormatException e)
                {
                   qMgrProp.Add(MQC.PORT_PROPERTY, 1414);
                }
    
                if (inParms.ContainsKey("-u"))
                   qMgrProp.Add(MQC.USER_ID_PROPERTY, ((System.String)inParms["-u"]));
    
                if (inParms.ContainsKey("-x"))
                   qMgrProp.Add(MQC.PASSWORD_PROPERTY, ((System.String)inParms["-x"]));
    
                System.Console.Out.WriteLine("MQTest62:");
                Console.WriteLine("  QMgrName ='{0}'", qManager);
                Console.WriteLine("  Output QName ='{0}'", inputQName);
    
                System.Console.Out.WriteLine("QMgr Property values:");
                foreach (DictionaryEntry de in qMgrProp)
                {
                   Console.WriteLine("  {0} = '{1}'", de.Key, de.Value);
                }
             }
             else
             {
                throw new System.ArgumentException();
             }
          }
    
          /// <summary> Connect, open queue, read (browse) a message, close queue and disconnect. </summary>
          ///
          private void testReceive()
          {
             MQQueueManager qMgr = null;
             MQQueue inQ = null;
             int openOptions = MQC.MQOO_INPUT_AS_Q_DEF + MQC.MQOO_FAIL_IF_QUIESCING;
    
             try
             {
                qMgr = new MQQueueManager(qManager, qMgrProp);
                System.Console.Out.WriteLine("MQTest62 successfully connected to " + qManager);
    
                inQ = qMgr.AccessQueue(inputQName, openOptions);
                System.Console.Out.WriteLine("MQTest62 successfully opened " + inputQName);
    
                testLoop(inQ);
    
             }
             catch (MQException mqex)
             {
                System.Console.Out.WriteLine("MQTest62 cc=" + mqex.CompletionCode + " : rc=" + mqex.ReasonCode);
             }
             catch (System.IO.IOException ioex)
             {
                System.Console.Out.WriteLine("MQTest62 ioex=" + ioex);
             }
             finally
             {
                try
                {
                   if (inQ != null)
                      inQ.Close();
                   System.Console.Out.WriteLine("MQTest62 closed: " + inputQName);
                }
                catch (MQException mqex)
                {
                   System.Console.Out.WriteLine("MQTest62 cc=" + mqex.CompletionCode + " : rc=" + mqex.ReasonCode);
                }
    
                try
                {
                   if (qMgr != null)
                      qMgr.Disconnect();
                   System.Console.Out.WriteLine("MQTest62 disconnected from " + qManager);
                }
                catch (MQException mqex)
                {
                   System.Console.Out.WriteLine("MQTest62 cc=" + mqex.CompletionCode + " : rc=" + mqex.ReasonCode);
                }
             }
          }
    
          private void testLoop(MQQueue inQ)
          {
             bool flag = true;
             MQGetMessageOptions gmo = new MQGetMessageOptions();
             gmo.Options |= MQC.MQGMO_WAIT | MQC.MQGMO_FAIL_IF_QUIESCING;
             gmo.WaitInterval = 2500;  // 2.5 seconds wait time or use MQC.MQEI_UNLIMITED to wait forever
             MQMessage msg = null;
    
             while (flag)
             {
                try
                {
                   msg = new MQMessage();
                   inQ.Get(msg, gmo);
                   System.Console.Out.WriteLine("Message Data: " + msg.ReadString(msg.MessageLength));
                }
                catch (MQException mqex)
                {
                   System.Console.Out.WriteLine("MQTest62 CC=" + mqex.CompletionCode + " : RC=" + mqex.ReasonCode);
                   if (mqex.Reason == MQC.MQRC_NO_MSG_AVAILABLE)
                   {
                      // no meesage - life is good - loop again
                   }
                   else
                   {
                      flag = false;  // severe error - time to exit
                   }
                }
                catch (System.IO.IOException ioex)
                {
                   System.Console.Out.WriteLine("MQTest62 ioex=" + ioex);
                }
             }
          }
    
          /// <summary> main line</summary>
          /// <param name="args">
          /// </param>
          //        [STAThread]
          public static void Main(System.String[] args)
          {
             MQTest62 write = new MQTest62();
    
             try
             {
                write.init(args);
                write.testReceive();
             }
             catch (System.ArgumentException e)
             {
                System.Console.Out.WriteLine("Usage: MQTest62 -h host -p port -c channel -m QueueManagerName -q QueueName [-u userID] [-x passwd]");
                System.Environment.Exit(1);
             }
             catch (MQException e)
             {
                System.Console.Out.WriteLine(e);
                System.Environment.Exit(1);
             }
    
             System.Environment.Exit(0);
          }
       }
    }
    

    【讨论】:

    • 很好的样品罗杰。感谢您的所有贡献。
    • @Roger 我错过了您的回复(现在我不得不再次管理代码),我还有另一个问题,如果出现 2009 错误消息 (MQRC_CONNECTION_BROKEN),是否可以在不退出应用程序的情况下重新连接?谢谢
    • 当然。关闭并断开连接,等待“x”秒然后连接,如果成功打开队列,则休眠“x”秒并重试。
    • 但是我必须创建一个新的队列管理器吗?我错过了一件非常简单的事情。我已经修改了您的代码以连接到我拥有的所有节点,我强制关闭连接(以模拟服务器关闭或网络问题),我收到了 2009 消息。我关闭连接。然后尝试再次打开,但我得到了 {"2018"}。我已指定重新连接,但似乎没有自动重新打开连接
    • 您是否查看了 MQ 知识中心中的信息?这是Java/MQ 的一个很好的解释,这与C# .NET 完全相同。另外,您为什么不查看 MQ 中包含的示例? SimpleClientAutoReconnectPut.cs 是一个执行自动重新连接的 C# .NET MQ 示例。
    猜你喜欢
    • 1970-01-01
    • 2011-03-25
    • 1970-01-01
    • 1970-01-01
    • 2016-01-28
    • 1970-01-01
    • 2015-05-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多