【发布时间】:2019-09-10 21:48:38
【问题描述】:
去年我开发了一个队列监视器,它使用 System.Reactive.Linq 检查 IBM MQ 总线上是否有消息
代码如下
public class QueueMonitor : IObservable<Message>, IDisposable
{
private string queueName;
private readonly MQQueue mqQueue;
private readonly MQQueueManager queueManager;
private readonly IDisposable timer;
private readonly object lockObj = new object();
private bool isChecking;
private readonly List<IObserver<Message>> observers;
public QueueMonitor(MQQueueManager queueManager, string queueName)
{
this.queueName = queueName;
this.queueManager = queueManager;
observers = new List<IObserver<Message>>();
mqQueue = queueManager.AccessQueue(queueName,
MQC.MQOO_INPUT_AS_Q_DEF // open queue for input
+ MQC.MQOO_FAIL_IF_QUIESCING); // but not if MQM stopping
timer = Observable.Interval(TimeSpan.FromSeconds(5)).Subscribe(_ =>
{
lock (lockObj)
{
if (!isChecking)
{
isChecking = true;
var mqMsg = new MQMessage();
var mqGetMsgOpts = new MQGetMessageOptions {WaitInterval = 1};
// 15 second limit for waiting
mqGetMsgOpts.Options |= MQC.MQGMO_WAIT;
try
{
mqQueue.Get(mqMsg, mqGetMsgOpts);
if (mqMsg.Format.CompareTo(MQC.MQFMT_STRING) == 0)
{
var text = mqMsg.ReadString(mqMsg.MessageLength);
System.Console.WriteLine(text);
Message message = new Message { Content = text };
foreach (var observer in observers)
observer.OnNext(message);
}
else
{
System.Console.WriteLine("Non-text message");
}
}
catch (MQException ex)
{
if ((ex.Message == "MQRC_NO_MSG_AVAILABLE"))
{
//nothing to do, emtpy queue
}
else
{
//log
}
}
finally
{
isChecking = false;
}
}
}
});
}
public IDisposable Subscribe(IObserver<Message> observer)
{
if (!observers.Contains(observer))
observers.Add(observer);
return new Unsubscriber(observers, observer);
}
public void Dispose()
{
((IDisposable)mqQueue)?.Dispose();
((IDisposable)queueManager)?.Dispose();
timer?.Dispose();
}
}
public class Unsubscriber : IDisposable
{
private readonly List<IObserver<Message>> _observers;
private readonly IObserver<Message> _observer;
public Unsubscriber(List<IObserver<Message>> observers, IObserver<Message> observer)
{
this._observers = observers;
this._observer = observer;
}
public void Dispose()
{
if (_observer != null) _observers.Remove(_observer);
}
}
这工作了将近一年,但现在有两件事需要修复,希望您能帮助我把它做好。
1)如果重启了IBMMQ,目前QueueMonitor没有收到新的传入消息,需要重启。
我应该如何处理?不知道Monitor端有没有重启IBM MQ。
2) 更复杂。我们正在迁移到一个新的平衡 IBMMQ 集群。它有 4 个配置为活动的活动节点。它们都在负载均衡器后面,所以当我在总线上放一条消息时,我将它发送到一个地址。
发送消息很简单。我遇到的问题是当我需要从队列中读取时。因为有 4 个不同的 IBMMQ 节点,有 4 个 IP。我怎么知道总线上已经发送了一条消息?我不能简单地听平衡器,因为它没有通知。我应该 ping 4 个节点吗?
平衡器是 netscaler。
提前致谢
【问题讨论】:
-
在捕获
MQException的逻辑中,当它不是MQRC_NO_MSG_AVAILABLE时,您是否会收到事件,但有一些其他错误表明MQ 正在重新启动或不可用? -
是的,我发现它会抛出 MQRC_CONNECTION_BROKEN,现在我检查是否可以执行重新连接。谢谢
-
我建议您不使用 LB 地址单独监控 4 个实例中的每一个。
-
@JoshMc 监视 4 个实例意味着在我的 queueMonitor 类中我必须注册这 4 个 IP 地址及其配置......我认为 c# 类具有相同的方法。你能告诉我你指的是哪个班级吗?
-
如果四个队列管理器同时都处于活动状态,您将不得不直接连接到它们中的每一个,因为您不会像您提到的那样“知道”消息被发送到哪个队列管理器。