重构消息处理框架ExchangeTask 只做新消息获取和分发

This commit is contained in:
zhouxiunai
2018-12-06 14:35:12 +08:00
parent 9687e735ee
commit 3e588fecf6
18 changed files with 225 additions and 183 deletions
@@ -14,5 +14,5 @@ public interface CminmsgDao extends CrudRepository<Cminmsg, Long> {
public List<Cminmsg> findByCminmsgsDateReceivedAfterOrderByCminmsgsDateReceived(Date date);
public List<Cminmsg> findByCminmsgsIdGreaterThanOrderByCminmsgsDateReceived(Long id);
public List<Cminmsg> findByCminmsgsIdGreaterThanAndCminmsgsDateProcessedIsNull(Long id);
}
@@ -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);
}
@@ -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
@@ -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 {
@@ -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 {
@@ -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 {
@@ -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 {
@@ -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 {
@@ -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));
}
}
@@ -17,7 +17,7 @@ import org.springframework.util.CollectionUtils;
* @author Administrator
*
*/
@Service
@Service("redisService")
public final class RedisService {
@Autowired
private RedisTemplate<String, Object> redisTemplate;
@@ -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<Cminmsg> getNewMsgsAfterId(Long id) {
List<Cminmsg> lsCminmsgs = cminmsgDao.findByCminmsgsIdGreaterThanOrderByCminmsgsDateReceived(id);
List<Cminmsg> lsCminmsgs = cminmsgDao.findByCminmsgsIdGreaterThanAndCminmsgsDateProcessedIsNull(id);
if(null == lsCminmsgs)
{
return Collections.emptyList();
@@ -57,18 +61,37 @@ public class CminmsgServiceImpl implements ICminmsgService {
@Override
public Optional<Cminmsg> getRespSchdCminmsg(Date date) {
List<Cminmsg> lsCminmsgs = getNewMsgsAfterDate(date);
Optional<Cminmsg> 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<Cminmsg> getRespSchdCminmsg(Date date,Integer tryTimes,Long intervalMs) {
Optional<Cminmsg> opCminmsg = Optional.empty();
Integer currentTimes=0;//当前尝试次数
while (currentTimes < tryTimes) {
currentTimes++;
List<Cminmsg> 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<Cminmsg> lsCmins = (List<Cminmsg>) cminmsgDao.findAll(ids);
List<Cminmsg> lsCminsProcessed = lsCmins.stream().map(node->{
node.setCminmsgsDateReceived(new Date());
node.setCminmsgsDateProcessed(new Date());
return node;
}).collect(Collectors.toList());
@@ -24,10 +24,12 @@ public interface ICminmsgService {
/**
* 获取回复的航班计划信息
* @param date
* @param date 大于当前日期
* @param 失败后尝试次数
* @param 每次尝试间隔毫秒数
* @return
*/
public Optional<Cminmsg> getRespSchdCminmsg(Date date);
public Optional<Cminmsg> getRespSchdCminmsg(Date date,Integer tryTimes,Long intervalMs);
/**
@@ -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
@@ -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<UpdateByCminmsgsResult> updateByCminmsgs(List<Cminmsg> lsCminmsgs) {
List<UpdateByCminmsgsResult> 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
@@ -19,13 +19,6 @@ public interface IFlightInfoService {
*/
public Boolean updateByDaySchd(Cminmsg cminmsg);
/**
* 根据消息更新航班信息
* @param lsCminmsgs
* @return 返回被更新了的动态航班信息
*/
public List<UpdateByCminmsgsResult> updateByCminmsgs(List<Cminmsg> lsCminmsgs);
/**
* 通过航班id获取航班动态
* @param flid
@@ -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<Cminmsg> opCminmsg = cminmsgService.getRespSchdCminmsg(waitSchdDate);
if(opCminmsg.isPresent())
{
Optional<Cminmsg> 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<UpdateByCminmsgsResult> 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<Cminmsg> 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