From 3e588fecf6b0a76257162e5e3ceda7a953a92654 Mon Sep 17 00:00:00 2001 From: zhouxiunai <154707516@qq.com> Date: Thu, 6 Dec 2018 14:35:12 +0800 Subject: [PATCH] =?UTF-8?q?=E9=87=8D=E6=9E=84=E6=B6=88=E6=81=AF=E5=A4=84?= =?UTF-8?q?=E7=90=86=E6=A1=86=E6=9E=B6ExchangeTask=20=E5=8F=AA=E5=81=9A?= =?UTF-8?q?=E6=96=B0=E6=B6=88=E6=81=AF=E8=8E=B7=E5=8F=96=E5=92=8C=E5=88=86?= =?UTF-8?q?=E5=8F=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../omms/msgexchangeapi/dao/CminmsgDao.java | 2 +- .../msghandler/IBaseHandler.java | 4 +- .../msghandler/MsgHandlerDispatcher.java | 10 +- .../msghandler/flop/ACFTHandler.java | 23 +++-- .../msghandler/flop/ACTTHandler.java | 20 ++-- .../msghandler/flop/BDPBHandler.java | 13 ++- .../msghandler/flop/BOTMHandler.java | 20 ++-- .../msghandler/flop/CKDTHandler.java | 20 ++-- .../msghandler/flop/FlopBaseHandler.java | 70 +++++++++++++ .../msgexchangeapi/redis/RedisService.java | 2 +- .../service/CminmsgServiceImpl.java | 53 +++++++--- .../service/ICminmsgService.java | 6 +- .../service/KafkaServiceImpl.java | 2 +- .../flightInfo/FlightInfoServiceImpl.java | 24 ----- .../flightInfo/IFlightInfoService.java | 7 -- .../msgexchangeapi/task/ExchangeTask.java | 98 +++---------------- src/main/resources/application-dev.yml | 2 + .../msgexchangeapi/dao/CminmsgDaoTest.java | 32 ++++++ 18 files changed, 225 insertions(+), 183 deletions(-) create mode 100644 src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/FlopBaseHandler.java create mode 100644 src/test/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDaoTest.java diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDao.java b/src/main/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDao.java index be8abea1..f2a4e2cd 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDao.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDao.java @@ -14,5 +14,5 @@ public interface CminmsgDao extends CrudRepository { public List findByCminmsgsDateReceivedAfterOrderByCminmsgsDateReceived(Date date); - public List findByCminmsgsIdGreaterThanOrderByCminmsgsDateReceived(Long id); + public List findByCminmsgsIdGreaterThanAndCminmsgsDateProcessedIsNull(Long id); } diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/IBaseHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/IBaseHandler.java index 0b09be17..05f0d502 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/IBaseHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/IBaseHandler.java @@ -1,5 +1,7 @@ package com.gzzn.omms.msgexchangeapi.msghandler; +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; + public interface IBaseHandler { - public HandlerResult run(String msg); + public HandlerResult run(Cminmsg cminmsg); } diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/MsgHandlerDispatcher.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/MsgHandlerDispatcher.java index 75d7f63a..fe91be6d 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/MsgHandlerDispatcher.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/MsgHandlerDispatcher.java @@ -8,9 +8,9 @@ import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.msg.Msg; import com.gzzn.omms.msgexchangeapi.service.IExchangeService; -import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; @Service public class MsgHandlerDispatcher { @@ -19,9 +19,9 @@ public class MsgHandlerDispatcher { @Autowired IExchangeService exchangeService; - public HandlerResult dispatch(String xmlMsg) + public HandlerResult dispatch(Cminmsg cminmsg) { - Msg msg = exchangeService.xmlstrToObject(xmlMsg, Msg.class); + Msg msg = exchangeService.xmlstrToObject(cminmsg.getCminmsgsClobMsg(), Msg.class); String type = msg.getMeta().getType(); String subType = msg.getMeta().getStyp(); @@ -29,8 +29,8 @@ public class MsgHandlerDispatcher { { Class clazzHandler = Class.forName(getHandlerClassName(type,subType)); Object classObject = clazzHandler.newInstance(); - Method runMethod = clazzHandler.getMethod("run", String.class); - HandlerResult result = (HandlerResult) runMethod.invoke(classObject, xmlMsg); + Method runMethod = clazzHandler.getMethod("run", Cminmsg.class); + HandlerResult result = (HandlerResult) runMethod.invoke(classObject, cminmsg); return result; } catch (ClassNotFoundException | InstantiationException diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/ACFTHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/ACFTHandler.java index 2d171de9..218e028f 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/ACFTHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/ACFTHandler.java @@ -1,22 +1,22 @@ package com.gzzn.omms.msgexchangeapi.msghandler.flop; +import java.util.Arrays; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.msg.flop.acft.ACFTMsg; import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.FLTR; import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; -import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler; -import com.gzzn.omms.msgexchangeapi.service.IExchangeService; -import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; -import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; +import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; /** * 航班机型事件 * @author zhouxiunai * */ -public class ACFTHandler implements IBaseHandler { +public class ACFTHandler extends FlopBaseHandler { private static Logger logger = LoggerFactory.getLogger(ACFTHandler.class); @@ -26,16 +26,19 @@ public class ACFTHandler implements IBaseHandler { * @return 处理结果及影响的动态航班信息 */ @Override - public HandlerResult run(String msg) { - IExchangeService exchangeService = (IExchangeService) SpringUtil.getBean("exchangeService"); - IFlightInfoService flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService"); - - ACFTMsg acftMsg = exchangeService.xmlstrToObject(msg, ACFTMsg.class); + public HandlerResult run(Cminmsg cminmsg) { + ACFTMsg acftMsg = exchangeService.xmlstrToObject(cminmsg.getCminmsgsClobMsg(), ACFTMsg.class); FLTR fltr = flightInfoService.getByFlid(acftMsg.getFLOP().getFLID()); if(null != fltr) { + //save fltr.setACFT(acftMsg.getFLOP().getACFT()); flightInfoService.saveFltr(fltr); + + //send and update + this.sendFltrAndMsg(fltr, cminmsg); + this.updateCminmsgToProcessed(cminmsg); + return HandlerResult.success(fltr); } else { diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/ACTTHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/ACTTHandler.java index a0967c51..11951901 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/ACTTHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/ACTTHandler.java @@ -3,33 +3,33 @@ package com.gzzn.omms.msgexchangeapi.msghandler.flop; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.msg.flop.actt.ACTTMsg; import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.FLTR; import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; -import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler; -import com.gzzn.omms.msgexchangeapi.service.IExchangeService; -import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; -import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; /** * 航班实际时间事件 * @author zhouxiunai * */ -public class ACTTHandler implements IBaseHandler{ +public class ACTTHandler extends FlopBaseHandler{ private static Logger logger = LoggerFactory.getLogger(ACTTHandler.class); @Override - public HandlerResult run(String msg) { - IExchangeService exchangeService = (IExchangeService) SpringUtil.getBean("exchangeService"); - IFlightInfoService flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService"); - - ACTTMsg acttMsg = exchangeService.xmlstrToObject(msg, ACTTMsg.class); + public HandlerResult run(Cminmsg cminmsg) { + ACTTMsg acttMsg = exchangeService.xmlstrToObject(cminmsg.getCminmsgsClobMsg(), ACTTMsg.class); FLTR fltr = flightInfoService.getByFlid(acttMsg.getFLOP().getFLID()); if(null != fltr) { + //save fltr.setACTT(acttMsg.getFLOP().getACTT()); flightInfoService.saveFltr(fltr); + + //send and update + this.sendFltrAndMsg(fltr, cminmsg); + this.updateCminmsgToProcessed(cminmsg); + return HandlerResult.success(fltr); } else { diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/BDPBHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/BDPBHandler.java index de859a25..f5948907 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/BDPBHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/BDPBHandler.java @@ -3,29 +3,34 @@ package com.gzzn.omms.msgexchangeapi.msghandler.flop; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.msg.flop.bdpb.BDPBMsg; import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.FLTR; import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; -import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler; import com.gzzn.omms.msgexchangeapi.service.IExchangeService; import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; -public class BDPBHandler implements IBaseHandler { +public class BDPBHandler extends FlopBaseHandler { private static Logger logger = LoggerFactory.getLogger(ACFTHandler.class); @Override - public HandlerResult run(String msg) { + public HandlerResult run(Cminmsg cminmsg) { IExchangeService exchangeService = (IExchangeService) SpringUtil.getBean("exchangeService"); IFlightInfoService flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService"); - BDPBMsg bdpbMsg = exchangeService.xmlstrToObject(msg, BDPBMsg.class); + BDPBMsg bdpbMsg = exchangeService.xmlstrToObject(cminmsg.getCminmsgsClobMsg(), BDPBMsg.class); FLTR fltr = flightInfoService.getByFlid(bdpbMsg.getFLOP().getFLID()); if(null != fltr) { //TO: 文档未描述该事件是什么 //fltr.setBOTM(bdpbMsg.getFLOP().getBDPB()); //flightInfoService.saveFltr(fltr); + + //send and update + this.sendFltrAndMsg(fltr, cminmsg); + this.updateCminmsgToProcessed(cminmsg); + return HandlerResult.success(fltr); } else { diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/BOTMHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/BOTMHandler.java index 73c8b3cb..b0cc6f4f 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/BOTMHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/BOTMHandler.java @@ -3,33 +3,33 @@ package com.gzzn.omms.msgexchangeapi.msghandler.flop; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.msg.flop.botm.BOTMMsg; import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.FLTR; import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; -import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler; -import com.gzzn.omms.msgexchangeapi.service.IExchangeService; -import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; -import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; /** * 航班登机时间事件处理 * @author zhouxiunai * */ -public class BOTMHandler implements IBaseHandler { +public class BOTMHandler extends FlopBaseHandler { private static Logger logger = LoggerFactory.getLogger(ACFTHandler.class); @Override - public HandlerResult run(String msg) { - IExchangeService exchangeService = (IExchangeService) SpringUtil.getBean("exchangeService"); - IFlightInfoService flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService"); - - BOTMMsg botmMsg = exchangeService.xmlstrToObject(msg, BOTMMsg.class); + public HandlerResult run(Cminmsg cminmsg) { + BOTMMsg botmMsg = exchangeService.xmlstrToObject(cminmsg.getCminmsgsClobMsg(), BOTMMsg.class); FLTR fltr = flightInfoService.getByFlid(botmMsg.getFLOP().getFLID()); if(null != fltr) { + //save fltr.setBOTM(botmMsg.getFLOP().getBOTM()); flightInfoService.saveFltr(fltr); + + //send and update + this.sendFltrAndMsg(fltr, cminmsg); + this.updateCminmsgToProcessed(cminmsg); + return HandlerResult.success(fltr); } else { diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/CKDTHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/CKDTHandler.java index bccacd94..b26968ff 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/CKDTHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/CKDTHandler.java @@ -3,33 +3,33 @@ package com.gzzn.omms.msgexchangeapi.msghandler.flop; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.msg.flop.ckdt.CKDTMsg; import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.FLTR; import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; -import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler; -import com.gzzn.omms.msgexchangeapi.service.IExchangeService; -import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; -import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; /** * 值机柜台事件处理 * @author zhouxiunai * */ -public class CKDTHandler implements IBaseHandler { +public class CKDTHandler extends FlopBaseHandler { private static Logger logger = LoggerFactory.getLogger(ACFTHandler.class); @Override - public HandlerResult run(String msg) { - IExchangeService exchangeService = (IExchangeService) SpringUtil.getBean("exchangeService"); - IFlightInfoService flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService"); - - CKDTMsg ckdtMsg = exchangeService.xmlstrToObject(msg, CKDTMsg.class); + public HandlerResult run(Cminmsg cminmsg) { + CKDTMsg ckdtMsg = exchangeService.xmlstrToObject(cminmsg.getCminmsgsClobMsg(), CKDTMsg.class); FLTR fltr = flightInfoService.getByFlid(ckdtMsg.getFLOP().getFLID()); if(null != fltr) { + //save fltr.setCKDT(ckdtMsg.getFLOP().getCKDT()); flightInfoService.saveFltr(fltr); + + //send and update + this.sendFltrAndMsg(fltr, cminmsg); + this.updateCminmsgToProcessed(cminmsg); + return HandlerResult.success(fltr); } else { diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/FlopBaseHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/FlopBaseHandler.java new file mode 100644 index 00000000..4baa3622 --- /dev/null +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/FlopBaseHandler.java @@ -0,0 +1,70 @@ +package com.gzzn.omms.msgexchangeapi.msghandler.flop; + +import java.util.Arrays; + +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; +import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.FLTR; +import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; +import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler; +import com.gzzn.omms.msgexchangeapi.redis.RedisService; +import com.gzzn.omms.msgexchangeapi.service.ICminmsgService; +import com.gzzn.omms.msgexchangeapi.service.IExchangeService; +import com.gzzn.omms.msgexchangeapi.service.IKafkaService; +import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; +import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; +import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; + +/** + * flop 航班动态消息事件的基类,代为注入了 exchangeService,flightInfoService + * @author Administrator + * + */ +public class FlopBaseHandler implements IBaseHandler { + protected IExchangeService exchangeService; + protected IFlightInfoService flightInfoService; + protected IKafkaService kafkaservice; + protected RedisService redisService; + protected ICminmsgService cminmsgService; + + FlopBaseHandler() + { + exchangeService = (IExchangeService) SpringUtil.getBean("exchangeService"); + flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService"); + kafkaservice = (IKafkaService) SpringUtil.getBean("kafkaservice"); + redisService = (RedisService) SpringUtil.getBean("redisService"); + cminmsgService = (ICminmsgService) SpringUtil.getBean("cminmsgService"); + } + + + @Override + public HandlerResult run(Cminmsg cminmsg) { + return null; + } + + /** + * 发送航班信息到schd topic ,同时发送动态消息到msg topic + * @param fltr + * @param cminmsg + */ + protected void sendFltrAndMsg(FLTR fltr,Cminmsg cminmsg) + { + //send schd + String schdMsg = JsonUtil.getString(fltr); + kafkaservice.msgSend("schd", schdMsg); + + //send msg + String msg = exchangeService.xmlToJson( + cminmsg.getCminmsgsClobMsg() + ); + kafkaservice.msgSend("msg", msg); + } + + /** + * 更新Cminmsg状态为已处理 + * @param cminmsg + */ + protected void updateCminmsgToProcessed(Cminmsg cminmsg) + { + cminmsgService.updateBatchProcessed(Arrays.asList(cminmsg)); + } +} diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/redis/RedisService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/redis/RedisService.java index 41071d02..88e14b1d 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/redis/RedisService.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/redis/RedisService.java @@ -17,7 +17,7 @@ import org.springframework.util.CollectionUtils; * @author Administrator * */ -@Service +@Service("redisService") public final class RedisService { @Autowired private RedisTemplate redisTemplate; diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/CminmsgServiceImpl.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/CminmsgServiceImpl.java index 26426e90..a478a5f6 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/CminmsgServiceImpl.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/CminmsgServiceImpl.java @@ -6,6 +6,8 @@ import java.util.List; import java.util.Optional; import java.util.stream.Collectors; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -13,8 +15,10 @@ import com.gzzn.omms.msgexchangeapi.dao.CminmsgDao; import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.msg.Msg; -@Service +@Service("cminmsgService") public class CminmsgServiceImpl implements ICminmsgService { + private static Logger logger = LoggerFactory.getLogger(CminmsgServiceImpl.class); + @Autowired private CminmsgDao cminmsgDao; @@ -46,7 +50,7 @@ public class CminmsgServiceImpl implements ICminmsgService { @Override public List getNewMsgsAfterId(Long id) { - List lsCminmsgs = cminmsgDao.findByCminmsgsIdGreaterThanOrderByCminmsgsDateReceived(id); + List lsCminmsgs = cminmsgDao.findByCminmsgsIdGreaterThanAndCminmsgsDateProcessedIsNull(id); if(null == lsCminmsgs) { return Collections.emptyList(); @@ -57,18 +61,37 @@ public class CminmsgServiceImpl implements ICminmsgService { @Override - public Optional getRespSchdCminmsg(Date date) { - List lsCminmsgs = getNewMsgsAfterDate(date); - Optional opCminmsg = lsCminmsgs - .stream() - .filter(x->{ - String msgBlob = x.getCminmsgsClobMsg(); - Msg msg = exchangeService.xmlstrToObject(msgBlob, Msg.class); - String type = msg.getMeta().getType(); - String subType = msg.getMeta().getStyp(); - return type.equals("SCHD") && subType.equals("RESP"); - }) - .findFirst(); + public Optional getRespSchdCminmsg(Date date,Integer tryTimes,Long intervalMs) { + Optional opCminmsg = Optional.empty(); + Integer currentTimes=0;//当前尝试次数 + + while (currentTimes < tryTimes) { + currentTimes++; + + List lsCminmsgs = getNewMsgsAfterDate(date); + opCminmsg = lsCminmsgs + .stream() + .filter(x->{ + String msgBlob = x.getCminmsgsClobMsg(); + Msg msg = exchangeService.xmlstrToObject(msgBlob, Msg.class); + String type = msg.getMeta().getType(); + String subType = msg.getMeta().getStyp(); + return type.equals("SCHD") && subType.equals("RESP"); + }) + .findFirst(); + // + if(opCminmsg.isPresent()) + { + break; + } + + //等待xx秒继续循环过程 + try { + Thread.sleep(intervalMs); + } catch (InterruptedException e) { + logger.warn("Thread.sleep:暂停线程异常"); + } + }//end while return opCminmsg; }//end function @@ -82,7 +105,7 @@ public class CminmsgServiceImpl implements ICminmsgService { List lsCmins = (List) cminmsgDao.findAll(ids); List lsCminsProcessed = lsCmins.stream().map(node->{ - node.setCminmsgsDateReceived(new Date()); + node.setCminmsgsDateProcessed(new Date()); return node; }).collect(Collectors.toList()); diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/ICminmsgService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/ICminmsgService.java index 7f984dd4..6736df86 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/ICminmsgService.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/ICminmsgService.java @@ -24,10 +24,12 @@ public interface ICminmsgService { /** * 获取回复的航班计划信息 - * @param date + * @param date 大于当前日期 + * @param 失败后尝试次数 + * @param 每次尝试间隔毫秒数 * @return */ - public Optional getRespSchdCminmsg(Date date); + public Optional getRespSchdCminmsg(Date date,Integer tryTimes,Long intervalMs); /** diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/KafkaServiceImpl.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/KafkaServiceImpl.java index 875c5508..646be0ff 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/KafkaServiceImpl.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/KafkaServiceImpl.java @@ -4,7 +4,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; -@Service +@Service("kafkaservice") public class KafkaServiceImpl implements IKafkaService{ @Autowired diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/flightInfo/FlightInfoServiceImpl.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/flightInfo/FlightInfoServiceImpl.java index 8e8fe427..df141c8f 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/flightInfo/FlightInfoServiceImpl.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/flightInfo/FlightInfoServiceImpl.java @@ -10,7 +10,6 @@ import org.springframework.stereotype.Service; import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.DNLDMsg; import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.FLTR; -import com.gzzn.omms.msgexchangeapi.msghandler.HandlerResult; import com.gzzn.omms.msgexchangeapi.msghandler.MsgHandlerDispatcher; import com.gzzn.omms.msgexchangeapi.redis.RedisService; import com.gzzn.omms.msgexchangeapi.service.IExchangeService; @@ -51,29 +50,6 @@ public class FlightInfoServiceImpl implements IFlightInfoService { return true; } //end function - - - @Override - public List updateByCminmsgs(List lsCminmsgs) { - List resultUpdated = new ArrayList();//更新结果 - - for (Cminmsg cminmsg : lsCminmsgs) { - //更新动态航班消息 - String strClobMsg = cminmsg.getCminmsgsClobMsg(); - HandlerResult handlerResult = msgHandlerDispatcher.dispatch(strClobMsg); - UpdateByCminmsgsResult updateResult = null; - if(handlerResult.getIsSuccess()) - { - updateResult = UpdateByCminmsgsResult.success(cminmsg, handlerResult.getFltrs()); - } - else { - updateResult = UpdateByCminmsgsResult.failure(cminmsg); - } - resultUpdated.add(updateResult); - }//end for - - return resultUpdated; - }//end function @Override diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/flightInfo/IFlightInfoService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/flightInfo/IFlightInfoService.java index 697c0137..17daec43 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/flightInfo/IFlightInfoService.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/flightInfo/IFlightInfoService.java @@ -19,13 +19,6 @@ public interface IFlightInfoService { */ public Boolean updateByDaySchd(Cminmsg cminmsg); - /** - * 根据消息更新航班信息 - * @param lsCminmsgs - * @return 返回被更新了的动态航班信息 - */ - public List updateByCminmsgs(List lsCminmsgs); - /** * 通过航班id获取航班动态 * @param flid diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java b/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java index cc830217..d9dfedbe 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java @@ -3,7 +3,6 @@ package com.gzzn.omms.msgexchangeapi.task; import java.util.Date; import java.util.List; import java.util.Optional; -import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -12,14 +11,11 @@ import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; +import com.gzzn.omms.msgexchangeapi.msghandler.MsgHandlerDispatcher; import com.gzzn.omms.msgexchangeapi.redis.RedisKeyConstant; import com.gzzn.omms.msgexchangeapi.redis.RedisService; import com.gzzn.omms.msgexchangeapi.service.ICminmsgService; -import com.gzzn.omms.msgexchangeapi.service.IExchangeService; -import com.gzzn.omms.msgexchangeapi.service.IKafkaService; import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; -import com.gzzn.omms.msgexchangeapi.service.flightInfo.UpdateByCminmsgsResult; -import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; /** @@ -31,14 +27,14 @@ import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; public class ExchangeTask { private static Logger logger = LoggerFactory.getLogger(ExchangeTask.class); - @Autowired - private IKafkaService kafkaservice; - + @Autowired + MsgHandlerDispatcher msgHandlerDispatcher; + + @Autowired RedisService redisService; - @Autowired - IExchangeService exchangeService; + @Autowired private ICminmsgService cminmsgService; @@ -53,27 +49,17 @@ public class ExchangeTask { logger.info("进入获取动态航班消息..."); Boolean isWaitSchd = (Boolean)redisService.get(RedisKeyConstant.KEY_ISWAITSCHD); - - - - Integer tryTimes = 0; - while (isWaitSchd && tryTimes < 10) { - logger.info("正在第{}次获取返回的动态航班日计划消息",tryTimes); - - tryTimes++; - - //查找回复的日计划消息 + if (isWaitSchd) { Date waitSchdDate = (Date)redisService.get(RedisKeyConstant.KEY_WAITSCHDDATE); - Optional opCminmsg = cminmsgService.getRespSchdCminmsg(waitSchdDate); - if(opCminmsg.isPresent()) - { + Optional opCminmsg = cminmsgService.getRespSchdCminmsg(waitSchdDate,10,3000L); + if(opCminmsg.isPresent()) { //如果找到,添加到内存数据库动态航班信息表 logger.info("找到返回的动态航班计划,开始更新动态航班信息..."); Boolean updateResult = flightInfoService.updateByDaySchd(opCminmsg.get()); if(false == updateResult) { logger.error("更新动态航班信息失败"); - break; + return ; } // @@ -86,25 +72,13 @@ public class ExchangeTask { redisService.set(RedisKeyConstant.KEY_ISWAITSCHD,isWaitSchd); logger.info("动态航班计划更新完成..."); - - break; //跳出循环 } else { - //等待xx秒继续循环过程 - try { - Thread.sleep(3000); - } catch (InterruptedException e) { - logger.warn("Thread.sleep:暂停线程异常"); - } + logger.error("获取动态航班消息失败,退出本次周期..."); + return ; } - } // end while 获取航班日记录 + } // end if 获取航班日记录 - if(isWaitSchd == true) - { - //获取返回的航班计划回复失败,发送报警信息退出整个程序 - logger.error("获取动态航班计划失败..."); - return ; - } //获取上次的beginId Integer beginIdInt = (Integer)redisService.get(RedisKeyConstant.KEY_LASTBEGINID); @@ -122,52 +96,12 @@ public class ExchangeTask { } logger.info("获取到{}条最新信息,准备更新动态航班信息",lsCminmsgs.size()); - List lsUpdatedFlightInfo = flightInfoService.updateByCminmsgs(lsCminmsgs); - //新增或变更 航班动态信息 发送 到kafka的SCHD topic - logger.info("开始发送schd信息..."); - for (UpdateByCminmsgsResult updateResult : lsUpdatedFlightInfo) { - if(updateResult.getIsSuccess()) - { - updateResult.getFltrs().forEach(node->{ - String jsonMsg = JsonUtil.getString(node); - kafkaservice.msgSend("schd", jsonMsg); - }); - } + //分发消息进行处理 + for (Cminmsg cminmsg : lsCminmsgs) { + msgHandlerDispatcher.dispatch(cminmsg); }//end for - //发送动态消息到kafka - logger.info("开始发送动态航班信息..."); - for (UpdateByCminmsgsResult updateResult : lsUpdatedFlightInfo) { - if(updateResult.getIsSuccess()) - { - String jsonMsg = exchangeService.xmlToJson( - updateResult - .getCminmsg() - .getCminmsgsClobMsg() - ); - - kafkaservice.msgSend("msg", jsonMsg); - }//end if - }//end for - - - //更新消息处理状态 - List lsCminmsgsSuccess = lsUpdatedFlightInfo.stream().filter(node->{ - return node.getIsSuccess(); - }).map(node->{ - return node.getCminmsg(); - }).collect(Collectors.toList()); - - - cminmsgService.updateBatchProcessed(lsCminmsgsSuccess); - - Long newBeginId = lsUpdatedFlightInfo - .get(lsUpdatedFlightInfo.size()-1) - .getCminmsg() - .getCminmsgsId(); - - redisService.set(RedisKeyConstant.KEY_LASTBEGINID,newBeginId); logger.info("本次定时获取新消息完成..."); } //end function diff --git a/src/main/resources/application-dev.yml b/src/main/resources/application-dev.yml index 11a85895..633fbd75 100644 --- a/src/main/resources/application-dev.yml +++ b/src/main/resources/application-dev.yml @@ -1,4 +1,6 @@ spring: + jpa: + show-sql: true datasource: username: ommsxc password: ommsxc diff --git a/src/test/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDaoTest.java b/src/test/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDaoTest.java new file mode 100644 index 00000000..046064cf --- /dev/null +++ b/src/test/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDaoTest.java @@ -0,0 +1,32 @@ +package com.gzzn.omms.msgexchangeapi.dao; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.List; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.junit4.SpringRunner; + +import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; +import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; + +@RunWith(SpringRunner.class) +@SpringBootTest +public class CminmsgDaoTest { + + @Autowired + private CminmsgDao cminmsgDao; + + @Test + public void testFindByCminmsgsIdGreaterThanAndCminmsgsDateProcessedIsNull() + { + List lsCminmsgs = cminmsgDao.findByCminmsgsIdGreaterThanAndCminmsgsDateProcessedIsNull(4501301L); + + assertThat(lsCminmsgs).isNotEmpty(); + + System.out.println(JsonUtil.getString(lsCminmsgs)); + } +}