From e920e676958536394e551a294f55292e29a47d22 Mon Sep 17 00:00:00 2001 From: zhouxiunai <154707516@qq.com> Date: Wed, 5 Dec 2018 09:52:51 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=8C=E6=88=90=E6=96=B0=E6=B6=88=E6=81=AF?= =?UTF-8?q?=E9=87=87=E9=9B=86=EF=BC=8C=E6=9B=B4=E6=96=B0=E5=8A=A8=E6=80=81?= =?UTF-8?q?=E8=88=AA=E7=8F=AD=E4=BF=A1=E6=81=AF=EF=BC=8C=E5=8F=91=E9=80=81?= =?UTF-8?q?=E5=8A=A8=E6=80=81=E8=88=AA=E7=8F=AD=E4=BF=A1=E6=81=AF=EF=BC=8C?= =?UTF-8?q?=E5=8F=91=E9=80=81=E5=8A=A8=E6=80=81=E8=88=AA=E7=8F=AD=E6=B6=88?= =?UTF-8?q?=E6=81=AF=E3=80=82=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../redis/RedisKeyConstant.java | 3 +- .../omms/msgexchangeapi/runner/AppRunner.java | 2 +- .../service/CminmsgServiceImpl.java | 18 ++++++ .../service/ICminmsgService.java | 7 +++ .../msgexchangeapi/task/ExchangeTask.java | 58 +++++++++++++++---- 5 files changed, 76 insertions(+), 12 deletions(-) diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/redis/RedisKeyConstant.java b/src/main/java/com/gzzn/omms/msgexchangeapi/redis/RedisKeyConstant.java index b61968c0..56ae3f21 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/redis/RedisKeyConstant.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/redis/RedisKeyConstant.java @@ -7,5 +7,6 @@ package com.gzzn.omms.msgexchangeapi.redis; */ public class RedisKeyConstant { public static String KEY_ISWAITSCHD="isWaitSchd"; - public static String KEY_MSGPROGRESS="msgProgress"; + public static String KEY_WAITSCHDDATE="waitSchdDate"; + public static String KEY_LASTBEGINID="lastBeginId"; } diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/runner/AppRunner.java b/src/main/java/com/gzzn/omms/msgexchangeapi/runner/AppRunner.java index 89cb1578..256d3039 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/runner/AppRunner.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/runner/AppRunner.java @@ -55,6 +55,6 @@ public class AppRunner implements CommandLineRunner { // redisService.set(RedisKeyConstant.KEY_ISWAITSCHD, true); //等待返回日计划状态 - redisService.set(RedisKeyConstant.KEY_MSGPROGRESS,nowDate); //初始化消息进度 + redisService.set(RedisKeyConstant.KEY_WAITSCHDDATE,nowDate); //初始化消息进度 }//end function run } 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 f4a3b7a0..26426e90 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/CminmsgServiceImpl.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/CminmsgServiceImpl.java @@ -4,6 +4,7 @@ import java.util.Collections; import java.util.Date; import java.util.List; import java.util.Optional; +import java.util.stream.Collectors; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -70,4 +71,21 @@ public class CminmsgServiceImpl implements ICminmsgService { .findFirst(); return opCminmsg; }//end function + + + @Override + public List updateBatchProcessed(List lsCminmsgs) { + // + List ids = lsCminmsgs.stream().map(node->{ + return node.getCminmsgsId(); + }).collect(Collectors.toList()); + + List lsCmins = (List) cminmsgDao.findAll(ids); + List lsCminsProcessed = lsCmins.stream().map(node->{ + node.setCminmsgsDateReceived(new Date()); + return node; + }).collect(Collectors.toList()); + + return (List) cminmsgDao.save(lsCminsProcessed); + }//end function } 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 8a3e2352..7f984dd4 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/ICminmsgService.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/ICminmsgService.java @@ -9,6 +9,13 @@ import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; public interface ICminmsgService { + /** + * 批量更新处理时间 + * @return + */ + public List updateBatchProcessed(List lsCminmsgs); + + /** * 发送xml格式的消息到cminmsgs表 * @param xmlMsg 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 761353bf..ccadd550 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java @@ -3,6 +3,7 @@ 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; @@ -10,7 +11,6 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; -import com.gzzn.omms.msgexchangeapi.entity.msg.schd.dnld.FLTR; import com.gzzn.omms.msgexchangeapi.redis.RedisKeyConstant; import com.gzzn.omms.msgexchangeapi.redis.RedisService; import com.gzzn.omms.msgexchangeapi.service.ICminmsgService; @@ -51,7 +51,7 @@ public class ExchangeTask { logger.info("定时任务启动...."); Boolean isWaitSchd = (Boolean)redisService.get(RedisKeyConstant.KEY_ISWAITSCHD); - Date msgProgress = (Date)redisService.get(RedisKeyConstant.KEY_MSGPROGRESS); + Long beginId = null; Integer tryTimes = 0; @@ -59,7 +59,8 @@ public class ExchangeTask { tryTimes++; //查找回复的日计划消息 - Optional opCminmsg = cminmsgService.getRespSchdCminmsg(msgProgress); + Date waitSchdDate = (Date)redisService.get(RedisKeyConstant.KEY_WAITSCHDDATE); + Optional opCminmsg = cminmsgService.getRespSchdCminmsg(waitSchdDate); if(opCminmsg.isPresent()) { //如果找到,添加到内存数据库动态航班信息表 @@ -84,29 +85,66 @@ public class ExchangeTask { } } // end while 获取航班日记录 - if(beginId == null) + if(isWaitSchd == true) { //获取返回的航班计划回复失败,发送报警信息退出整个程序 + logger.error("获取动态航班计划失败..."); return ; } + if(null == beginId) + { + //获取上次的beginId + beginId = (Long)redisService.get(RedisKeyConstant.KEY_LASTBEGINID); + } + //获取航班动态消息 List lsCminmsgs = cminmsgService.getNewMsgsAfterId(beginId); List lsUpdatedFlightInfo = flightInfoService.updateByCminmsgs(lsCminmsgs); //新增或变更 航班动态信息 发送 到kafka的SCHD topic for (UpdateByCminmsgsResult updateResult : lsUpdatedFlightInfo) { - - } + if(updateResult.getIsSuccess()) + { + updateResult.getFltrs().forEach(node->{ + String jsonMsg = exchangeService.objectToXmlstr(node); + kafkaservice.msgSend("schd", jsonMsg); + }); + } + }//end for //发送动态消息到kafka - for (Cminmsg cminmsg : lsCminmsgs) { - - } + 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 }