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

如何用ibmmq与TypeScript实现同步MQ调用并读取响应?

问题描述

我是TypeScript新手,有一定MQ知识,因为实现的方案一年后会弃用,所以够用即可的方案就满足需求。目前已经通过ibmmq成功向队列发送消息,但不知道怎么实现同步MQ调用并获取、读取响应。

当前可运行的代码如下:

mq.ConnxPromise(qMgr, cno)
    .then((hConn) => {
        logger.info("a) MQCONN to %s successful ", qMgr);
        ghConn = hConn;
        const od = new mq.MQOD();
        od.ObjectName = qName;
        od.ObjectType = MQC.MQOT_Q;
        const openOptions = MQC.MQOO_OUTPUT;
        return mq.OpenPromise(hConn, od, openOptions);
    })
    .then((hObj) => {
        logger.info("b) MQOPEN of %s successful", qName);

        const mqmd = new mq.MQMD(); // Defaults are fine.
        const pmo = new mq.MQPMO();
        // Describe how the Put should behave
        pmo.Options = MQC.MQPMO_NO_SYNCPOINT | MQC.MQPMO_NEW_MSG_ID | MQC.MQPMO_NEW_CORREL_ID;

        ghObj = hObj;
        const response = mq.PutPromise(hObj, mqmd, pmo, message);
        return response;
    })
    .then(() => {
        logger.info("MQPUT successful");
        return mq.ClosePromise(ghObj, 0);
    })
    .then(() => {
        logger.info("MQCLOSE successful");
        return mq.DiscPromise(ghConn);
    })
    .then(() => {
        logger.info("Done.");
    })
同步MQ调用获取响应的实现方案

因为是“够用即可”的需求,我们可以基于现有代码,通过**关联ID(CorrelID)**匹配请求和响应,核心逻辑是:给请求打一个唯一“标签”,发送后从响应队列里找带相同“标签”的消息。

步骤1:准备响应队列

先确认有可读取的响应队列(建议和请求队列分开),确保MQ配置允许你从该队列读取消息。

步骤2:修改发送逻辑,手动设置关联ID

发送请求时,不要让MQ自动生成关联ID,而是自己生成一个唯一ID作为标识,方便后续匹配响应。修改原代码的Put部分:

// 用uuid生成唯一ID(先安装依赖:npm install uuid,导入:import { v4 as UUID } from 'uuid')
const correlId = Buffer.from(UUID.v4().replace(/-/g, ""), "hex");

const mqmd = new mq.MQMD();
mqmd.CorrelId = correlId; // 手动设置关联ID

const pmo = new mq.MQPMO();
// 去掉MQPMO_NEW_CORREL_ID,因为我们自己指定了CorrelId
pmo.Options = MQC.MQPMO_NO_SYNCPOINT | MQC.MQPMO_NEW_MSG_ID;

步骤3:发送请求后读取响应

发送完请求后,不立即关闭连接,而是打开响应队列,循环读取消息直到找到匹配关联ID的响应。

完整修改代码(用async/await更易读)

把链式then改成async/await,逻辑更直观:

async function sendAndReceiveMessage(qMgr: string, cno: any, reqQName: string, respQName: string, message: string) {
    let hConn: any;
    let reqHObj: any;
    let respHObj: any;

    try {
        // 1. 连接MQ管理器
        hConn = await mq.ConnxPromise(qMgr, cno);
        logger.info("a) MQCONN to %s successful ", qMgr);

        // 2. 打开请求队列发送消息
        const reqOd = new mq.MQOD();
        reqOd.ObjectName = reqQName;
        reqOd.ObjectType = MQC.MQOT_Q;
        reqHObj = await mq.OpenPromise(hConn, reqOd, MQC.MQOO_OUTPUT);
        logger.info("b) MQOPEN of %s successful", reqQName);

        // 3. 生成唯一关联ID
        const correlId = Buffer.from(UUID.v4().replace(/-/g, ""), "hex");
        logger.info("Generated CorrelID: %s", correlId.toString("hex"));

        // 4. 发送请求消息
        const mqmd = new mq.MQMD();
        mqmd.CorrelId = correlId;
        const pmo = new mq.MQPMO();
        pmo.Options = MQC.MQPMO_NO_SYNCPOINT | MQC.MQPMO_NEW_MSG_ID;

        await mq.PutPromise(reqHObj, mqmd, pmo, message);
        logger.info("MQPUT successful");

        // 5. 关闭请求队列
        await mq.ClosePromise(reqHObj, 0);
        logger.info("MQCLOSE of request queue successful");

        // 6. 打开响应队列准备读取
        const respOd = new mq.MQOD();
        respOd.ObjectName = respQName;
        respOd.ObjectType = MQC.MQOT_Q;
        respHObj = await mq.OpenPromise(hConn, respOd, MQC.MQOO_INPUT_SHARED);
        logger.info("c) MQOPEN of %s successful", respQName);

        // 7. 循环读取响应,直到找到匹配的CorrelID或超时
        const gmo = new mq.MQGMO();
        gmo.WaitInterval = 30000; // 超时30秒
        gmo.Options = MQC.MQGMO_NO_SYNCPOINT | MQC.MQGMO_WAIT | MQC.MQGMO_CONVERT;

        let responseMessage: string | null = null;
        while (!responseMessage) {
            try {
                const getResponse = await mq.GetPromise(respHObj, gmo);
                const msgMd = getResponse.md;
                const msgBuffer = getResponse.message;

                // 匹配关联ID
                if (msgMd.CorrelId.equals(correlId)) {
                    responseMessage = msgBuffer.toString();
                    logger.info("Received matching response: %s", responseMessage);
                } else {
                    logger.info("Skipping message with mismatched CorrelID");
                }
            } catch (err: any) {
                if (err.mqrc === MQC.MQRC_NO_MSG_AVAILABLE) {
                    logger.info("No response received within timeout");
                    break;
                } else {
                    throw err;
                }
            }
        }

        return responseMessage;
    } catch (err) {
        logger.error("Error occurred: %s", err);
        throw err;
    } finally {
        // 确保资源被清理
        if (respHObj) await mq.ClosePromise(respHObj, 0).catch(() => {});
        if (hConn) await mq.DiscPromise(hConn).catch(() => {});
        logger.info("Done.");
    }
}

关键说明

  • 关联ID(CorrelID):相当于请求的“身份证”,服务端处理完请求后,需要把这个ID附在响应消息上,你才能找到对应的响应。
  • async/await:比链式then更像同步代码,新手更容易理解和调试。
  • 超时设置:避免无限等待,这里设了30秒,可根据需求调整。
  • 资源清理:finally块确保队列和连接一定会关闭,防止MQ资源泄漏。

额外注意点

  1. 要确保服务端逻辑会把请求的CorrelID复制到响应消息的CorrelID字段,否则无法匹配响应。
  2. 如果没有uuid库,也可以用简单的方式生成唯一ID(比如时间戳+随机字符串),只要保证每次请求唯一即可。
  3. 这个方案只满足基础同步需求,没有做复杂的重试、并发处理,符合“够用即可”的要求。

内容的提问来源于stack exchange,提问作者Richard Peters

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 12:04:53