From afa1c8d90da8b442cfaea8532fe305c9133ada9d Mon Sep 17 00:00:00 2001 From: zhouxiunai <154707516@qq.com> Date: Tue, 11 Dec 2018 13:59:59 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E8=AE=A1=E5=88=92=E5=A4=96?= =?UTF-8?q?=E8=88=AA=E7=8F=AD=E6=B6=88=E6=81=AF=E5=A4=84=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../msghandler/schd/ADFTHandler.java | 38 +++++++++++++++++++ .../msghandler/schd/base/SchdBaseHandler.java | 21 ++++++++++ 2 files changed, 59 insertions(+) create mode 100644 src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/ADFTHandler.java diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/ADFTHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/ADFTHandler.java new file mode 100644 index 00000000..f514b40e --- /dev/null +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/ADFTHandler.java @@ -0,0 +1,38 @@ +package com.gzzn.omms.msgexchangeapi.msghandler.schd; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; +import com.gzzn.omms.msgexchangeapi.entity.msg.MSG; +import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR; +import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; +import com.gzzn.omms.msgexchangeapi.msghandler.schd.base.SchdBaseHandler; + +/** + * 添加计划外航班计划 + * @author zhouxiunai + * + */ +public class ADFTHandler extends SchdBaseHandler{ + private static Logger logger = LoggerFactory.getLogger(ADFTHandler.class); + + @Override + public HandlerResult run(Cminmsg cminmsg) { + MSG msg = exchangeService.xmlstrToObject(cminmsg.getCminmsgsClobMsg(), MSG.class); + if(msg.getSCHD().getFLTR().size() != 1) + { + logger.error("预期每条消息只会包含一个动态航班,当前{}",msg.getSCHD().getFLTR().size()); + return HandlerResult.failure(); + } + + //每条消息只会包含一个动态航班 + FLTR fltr = msg.getSCHD().getFLTR().get(0); + flightInfoService.saveFltr(fltr); + + this.sendFltrAndMsg(fltr, cminmsg); + this.updateCminmsgToProcessed(cminmsg); + + return HandlerResult.success(); + } //end function +} diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/base/SchdBaseHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/base/SchdBaseHandler.java index daa17a93..b89b438c 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/base/SchdBaseHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/base/SchdBaseHandler.java @@ -3,12 +3,14 @@ package com.gzzn.omms.msgexchangeapi.msghandler.schd.base; import java.util.Arrays; import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; +import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD; import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler; import com.gzzn.omms.msgexchangeapi.service.ICminmsgService; import com.gzzn.omms.msgexchangeapi.service.IKafkaService; import com.gzzn.omms.msgexchangeapi.service.exchange.IExchangeService; import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; +import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; /** @@ -47,6 +49,25 @@ public class SchdBaseHandler implements IBaseHandler{ kafkaservice.msgSend("msg", msg); } + + /** + * 发送航班信息到schd topic ,同时发送动态消息到msg topic + * @param fltr + * @param cminmsg + */ + protected void sendFltrAndMsg(SCHD.FLTR fltr,Cminmsg cminmsg) + { + //send schd + String schdMsg = JsonUtil.getString(fltr); + kafkaservice.msgSend("schd", schdMsg); + + //send msg + String msg = exchangeService.xmlstrToJson( + cminmsg.getCminmsgsClobMsg() + ); + kafkaservice.msgSend("msg", msg); + } + /** * 更新Cminmsg状态为已处理 * @param cminmsg