如何并行发送ISOMsg请求至双渠道并对比结果?
问题分析与解决方案
当前实现的问题
你的串行实现不仅响应时间会叠加两个渠道的耗时,还存在几处逻辑和健壮性问题:
- 逻辑错误:当
host1返回成功后,你判断isoMsg.hasField(39)是不合理的——应该检查的是reply_bank是否包含字段39;而且最后你修改了reply的39字段为00,却给客户端发送reply_bank,这里逻辑矛盾。 - 异常处理不严谨:
ISOException被空catch,会隐藏字段解析、打包等关键错误,导致问题难以排查。 - 超时/无响应处理缺失:如果
mux.request()返回null(超时),当前代码直接跳过,客户端会一直等待回复,引发超时问题。 - 反转逻辑未实现:当
host2失败时,仅注释了要给host1发反转,没有具体实现,容易引发单边账风险。
并行处理实现方案
要实现并行请求两个渠道,我们可以用Java的ExecutorService线程池来异步发送请求,通过Future来获取两个渠道的响应,等两者都返回后再做结果对比和客户端回复。
优化后的代码示例
@Override public boolean process(ISOSource isoSrc, ISOMsg isoMsg) { // 只处理0200交易 if (!"0200".equals(isoMsg.getMTI())) { return false; } // 创建固定大小的线程池,专门处理两个渠道的并行请求 ExecutorService executor = Executors.newFixedThreadPool(2); try { // 复制请求报文,避免多线程修改同一对象引发线程安全问题 ISOMsg msgForHost1 = (ISOMsg) isoMsg.clone(); ISOMsg msgForHost2 = (ISOMsg) isoMsg.clone(); // 提交两个异步任务 Future<ISOMsg> host1Future = executor.submit(() -> { MUX mux = (MUX) NameRegistrar.getIfExists("mux.jpos-host1-mux"); return mux.request(msgForHost1, 10 * 1000); }); Future<ISOMsg> host2Future = executor.submit(() -> { MUX muxBank = (MUX) NameRegistrar.getIfExists("mux.jpos-host2-mux"); return muxBank.request(msgForHost2, 10 * 1000); }); // 等待两个任务完成,获取结果 ISOMsg replyHost1 = host1Future.get(11 * 1000, TimeUnit.MILLISECONDS); // 总超时略大于单个请求超时 ISOMsg replyHost2 = host2Future.get(11 * 1000, TimeUnit.MILLISECONDS); // 结果处理逻辑 ISOMsg clientReply = null; boolean host1Success = replyHost1 != null && "00".equals(replyHost1.getValue(39)); boolean host2Success = replyHost2 != null && "00".equals(replyHost2.getValue(39)); if (host1Success && host2Success) { // 两者都成功,构造成功回复给客户端(可根据需求选择用host1或host2的报文,或合并字段) clientReply = (ISOMsg) replyHost1.clone(); clientReply.set(39, "00"); } else if (host1Success && !host2Success) { // host1成功但host2失败,需要给host1发送反转交易 sendReversalToHost1(msgForHost1); // 给客户端返回失败 clientReply = (ISOMsg) replyHost1.clone(); clientReply.set(39, "05"); // 示例失败码,根据实际调整 } else { // host1失败,直接返回host1的失败回复给客户端 clientReply = replyHost1 != null ? (ISOMsg) replyHost1.clone() : buildErrorReply(isoMsg, "96"); // 系统错误码 } // 发送回复给客户端 if (clientReply != null) { isoSrc.send(clientReply); } } catch (ISOException e) { Logger.getLogger(ISO_HAWK.class.getName()).log(Level.SEVERE, "ISO报文处理异常", e); sendErrorReply(isoSrc, isoMsg, "96"); } catch (IOException e) { Logger.getLogger(ISO_HAWK.class.getName()).log(Level.SEVERE, "网络IO异常", e); sendErrorReply(isoSrc, isoMsg, "96"); } catch (InterruptedException | ExecutionException | TimeoutException e) { Logger.getLogger(ISO_HAWK.class.getName()).log(Level.SEVERE, "并行请求超时或中断", e); // 超时情况下,需要检查已完成的请求,必要时发送反转 sendErrorReply(isoSrc, isoMsg, "08"); // 超时错误码 } finally { // 关闭线程池 executor.shutdown(); } return false; } // 辅助方法:发送反转交易到host1 private void sendReversalToHost1(ISOMsg originalMsg) throws ISOException, IOException { ISOMsg reversalMsg = (ISOMsg) originalMsg.clone(); reversalMsg.setMTI("0400"); // 反转交易MTI,根据实际协议调整 // 设置其他反转所需字段,比如原交易的系统跟踪号等 reversalMsg.set(11, originalMsg.getValue(11)); MUX mux = (MUX) NameRegistrar.getIfExists("mux.jpos-host1-mux"); mux.request(reversalMsg, 10 * 1000); } // 辅助方法:构建错误回复 private void sendErrorReply(ISOSource isoSrc, ISOMsg originalMsg, String respCode) throws ISOException, IOException { ISOMsg errorReply = (ISOMsg) originalMsg.clone(); errorReply.setMTI("0210"); errorReply.set(39, respCode); isoSrc.send(errorReply); }
关键优化点
- 并行异步请求:通过
ExecutorService同时向两个渠道发送请求,总响应时间取决于耗时较长的那个渠道,而非两者之和。 - 线程安全:复制原始
ISOMsg对象,避免多线程修改同一实例引发的线程安全问题。 - 超时控制:通过
Future.get()设置总超时时间,防止无限等待。 - 健壮性增强:完善异常处理,补充反转逻辑和错误回复,避免单边账和客户端无响应问题。
- 清晰的结果判断:明确两个渠道的成功/失败组合场景,分别处理。
内容的提问来源于stack exchange,提问作者Dexter
相关产品推荐
相关产品推荐

