【发布时间】:2021-07-03 02:06:51
【问题描述】:
我正在使用低级直接 IBM MQ 库将消息推送到队列中并检索它们。我试图设置应用程序以便消息可以进来,比如通过提取数据库的数据然后将记录推送到队列中,实际上相同的代码可以读取消息。我主要是想设置一个线程,一旦消息出现就会拉取消息。
此代码运行,第一个 PUT 工作,但第二个不工作并挂起。我不明白这里的流程吗
另外,如果我从下面的第二个“GET”周围获取代码,我是否可以编写一个线程,每 500 毫秒调用一次该例程,等待新消息进入。
final int putOptions = MQC.MQPMO_NO_SYNCPOINT
| MQC.MQPMO_SYNC_RESPONSE;
this.mqPMO = new MQPutMessageOptions();
this.mqPMO.options = putOptions;
// This code hangs !!!! (error here)
mqueue.put(msg, this.mqPMO);
...
public void bootstap() {
MQEnvironment.hostname = "localhost";
MQEnvironment.port = 1414;
MQEnvironment.channel = "DEV.ADMIN.SVRCONN";
MQEnvironment.properties.put(MQConstants.APPNAME_PROPERTY, "my_application_name");
MQEnvironment.enableTracing(5);
MQQueueManager mqManager = null;
MQQueue mqueue = null;
try {
// MQCNO_CLIENT_BINDING is not available for Java or .NET as they have their own mechanisms for choosing the bind type.
final String qmName = "QM1";
final String userId = "admin";
final String Password = "passw0rd";
final Hashtable h = new Hashtable();
h.put(MQConstants.USER_ID_PROPERTY, userId);
h.put(MQConstants.PASSWORD_PROPERTY, Password);
h.put(MQConstants.USE_MQCSP_AUTHENTICATION_PROPERTY, true);
mqManager = new MQQueueManager(qmName, h);
//mqManager = new MQQueueManager(qmName, WMQConstants.WMQ_CM_BINDINGS);
this.mqGMO = new MQGetMessageOptions();
this.mqGMO.options = MQC.MQGMO_NO_SYNCPOINT |
MQC.MQGMO_WAIT |
MQC.MQGMO_CONVERT |
MQC.MQGMO_FAIL_IF_QUIESCING;
this.mqGMO.matchOptions = MQC.MQMO_MATCH_CORREL_ID;
this.mqGMO.waitInterval = MQC.MQWI_UNLIMITED;
int openOptions = MQC.MQOO_INPUT_SHARED |
MQC.MQOO_OUTPUT;
mqueue = mqManager.accessQueue("DEV.QUEUE.1", openOptions);
logger.info(">> Find connection handle queue manager - " + mqueue);
{
final MQMessage msg = new MQMessage();
final String correlId = "0002";
final String byteArry = this.hexStringToByteArray(correlId);
logger.info(">>> correlId: " + correlId);
logger.info(">>> byteArry: " + byteArry);
msg.correlationId = byteArry.getBytes();
msg.format = MQConstants.MQFMT_STRING;
// ... and write some text in UTF8 format
msg.writeUTF("{{ Hello, World }}}");
// Use the default put message options...
// Or: pmo.options = MQConstants.MQPMO_ASYNC_RESPONSE
final int putOptions = MQC.MQPMO_NO_SYNCPOINT
| MQC.MQPMO_SYNC_RESPONSE;
this.mqPMO = new MQPutMessageOptions();
this.mqPMO.options = putOptions;
// put the message //
mqueue.put(msg, this.mqPMO);
logger.info(" >>> Continue to get routine");
}
{
// This code works !!! get the message
MQMessage retrievedMessage = new MQMessage();
retrievedMessage.correlationId = this.hexStringToByteArray("0001").getBytes();
mqueue.get(retrievedMessage, this.mqGMO);
// And prove we have the message by displaying the UTF message text
String msgText = retrievedMessage.readUTF();
logger.info("~~~~ The message is: " + msgText);
}
{
final MQMessage msg = new MQMessage();
final String correlId = "0001";
final String byteArry = this.hexStringToByteArray(correlId);
logger.info(">>> correlId: " + correlId);
logger.info(">>> byteArry: " + byteArry);
msg.correlationId = byteArry.getBytes();
msg.format = MQConstants.MQFMT_STRING;
// ... and write some text in UTF8 format
msg.writeUTF("{{ Hello, World }}}");
// Use the default put message options...
// Or: pmo.options = MQConstants.MQPMO_ASYNC_RESPONSE
final int putOptions = MQC.MQPMO_NO_SYNCPOINT
| MQC.MQPMO_SYNC_RESPONSE;
this.mqPMO = new MQPutMessageOptions();
this.mqPMO.options = putOptions;
// This code hangs !!!! (error here)
mqueue.put(msg, this.mqPMO);
}
mqueue.close();
mqManager.disconnect();
} catch(final Exception e) {
logger.error("Error at MQ manager", e);
}
}
【问题讨论】:
-
您当然可以将 GET 代码移动到另一个线程。您可以为 waitInterval 指定所需的超时并让 GET 调用等待该间隔的消息,而不是等待 MQC.MQWI_UNLIMITED 消息。 PUT 挂起很有趣,你在做别的事情吗?完整的代码可能会有所帮助。
-
您使用的是哪个版本的 IBM MQ jar 文件?
-
IBM MQ 所有客户端:9.2.2.0 我确实尝试更改等待,这可能已修复它。现在我收到队列已清空的错误。也许我会走那条路。