diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/FDELHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/FDELHandler.java index 9572557a..6690cfab 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/FDELHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/FDELHandler.java @@ -52,7 +52,7 @@ public class FDELHandler extends FlopBaseHandler { }//end for //推送变更后的主航班信息 - this.sendFltr(fltrMafl); + this.addToBuffer(fltrMafl); } return HandlerResult.success(fltr); diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/base/FlopBaseHandler.java b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/base/FlopBaseHandler.java index 8ced18c2..3b8329f5 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/base/FlopBaseHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/flop/base/FlopBaseHandler.java @@ -15,9 +15,9 @@ 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.IKafkaService; +import com.gzzn.omms.msgexchangeapi.service.IMsgBufferService; 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; /** @@ -33,6 +33,7 @@ public class FlopBaseHandler implements IBaseHandler { protected IKafkaService kafkaservice; protected RedisService redisService; protected ICminmsgService cminmsgService; + protected IMsgBufferService msgbufferservice; public FlopBaseHandler() { @@ -41,6 +42,7 @@ public class FlopBaseHandler implements IBaseHandler { kafkaservice = (IKafkaService) SpringUtil.getBean("kafkaservice"); redisService = (RedisService) SpringUtil.getBean("redisService"); cminmsgService = (ICminmsgService) SpringUtil.getBean("cminmsgService"); + msgbufferservice = (IMsgBufferService) SpringUtil.getBean("msgbufferservice"); } @@ -60,7 +62,8 @@ public class FlopBaseHandler implements IBaseHandler { flightInfoService.saveFltr(fltrResult); //send and update - this.sendFltrAndMsg(fltrResult, cminmsg); + this.addToBuffer(fltrResult); + this.sendMsg(cminmsg); this.updateCminmsgToProcessed(cminmsg); return HandlerResult.success(fltrResult); @@ -73,6 +76,16 @@ public class FlopBaseHandler implements IBaseHandler { } } + /** + * 添加航班动态信息到缓冲区 + * @param fltr + */ + protected void addToBuffer(SCHD.FLTR fltr) + { + msgbufferservice.addFltr(fltr);//缓存区 + } + + /** * 关键处理逻辑 ,更新航班动态信息 * @param flop 消息 @@ -85,34 +98,6 @@ public class FlopBaseHandler implements IBaseHandler { } - /** - * 发送航班信息到schd topic ,同时发送动态消息到msg topic - * @param fltr - * @param cminmsg - */ - protected void sendFltrAndMsg(SCHD.FLTR fltr,CminmsgWapper cminmsg) - { - //send schd - String schdMsg = JsonUtil.getString(fltr); - kafkaservice.msgSend("schd", schdMsg); - - //send msg - String msg = exchangeService.msgToJson(cminmsg.getMsg()); - kafkaservice.msgSend("msg", msg); - } - - /** - * 发送动态航班信息 - * @param fltr - */ - protected void sendFltr(SCHD.FLTR fltr) - { - //send schd - String schdMsg = JsonUtil.getString(fltr); - kafkaservice.msgSend("schd", schdMsg); - } - - /** * 仅发送动态消息 * @param cminmsg 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 index 57f6a142..ea0a6c6e 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/ADFTHandler.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/msghandler/schd/ADFTHandler.java @@ -34,7 +34,8 @@ public class ADFTHandler extends SchdBaseHandler{ flightInfoService.saveFltr(fltr); - this.sendFltrAndMsg(fltr, cminmsg); + this.addToBuffer(fltr); + this.sendMsg(cminmsg); this.updateCminmsgToProcessed(cminmsg); @@ -48,7 +49,7 @@ public class ADFTHandler extends SchdBaseHandler{ fltrMafl.getMafl().add(mafldata); //推送变更后的主航班信息 - this.sendFltr(fltrMafl); + this.addToBuffer(fltrMafl); } return HandlerResult.success(); 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 28e65ec0..c8c60749 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 @@ -5,7 +5,6 @@ import java.util.Arrays; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.gzzn.omms.msgexchangeapi.dto.ResponseDto; import com.gzzn.omms.msgexchangeapi.entity.Cminmsg; import com.gzzn.omms.msgexchangeapi.entity.CminmsgWapper; import com.gzzn.omms.msgexchangeapi.entity.msg.MSG; @@ -14,9 +13,9 @@ 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.IMsgBufferService; 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; /** @@ -31,6 +30,7 @@ public class SchdBaseHandler implements IBaseHandler{ protected IFlightInfoService flightInfoService; protected IKafkaService kafkaservice; protected ICminmsgService cminmsgService; + protected IMsgBufferService msgbufferservice; public SchdBaseHandler() { @@ -38,6 +38,7 @@ public class SchdBaseHandler implements IBaseHandler{ flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService"); kafkaservice = (IKafkaService) SpringUtil.getBean("kafkaservice"); cminmsgService = (ICminmsgService) SpringUtil.getBean("cminmsgService"); + msgbufferservice = (IMsgBufferService) SpringUtil.getBean("msgbufferservice"); } @Override @@ -45,6 +46,15 @@ public class SchdBaseHandler implements IBaseHandler{ return null; } + /** + * 添加航班动态信息到缓冲区 + * @param fltr + */ + protected void addToBuffer(SCHD.FLTR fltr) + { + msgbufferservice.addFltr(fltr);//缓存区 + } + /** * 发送日计划到kafka队列 * @param cminmsg @@ -80,34 +90,6 @@ public class SchdBaseHandler implements IBaseHandler{ String msg = exchangeService.msgToJson(cminmsg.getMsg()); kafkaservice.msgSend("msg", msg); } - - - /** - * 发送航班信息到schd topic ,同时发送动态消息到msg topic - * @param fltr - * @param cminmsg - */ - protected void sendFltrAndMsg(SCHD.FLTR fltr,CminmsgWapper cminmsg) - { - //send schd - String schdMsg = JsonUtil.getString(fltr); - kafkaservice.msgSend("schd", schdMsg); - - //send msg - String msg = exchangeService.msgToJson(cminmsg.getMsg()); - kafkaservice.msgSend("msg", msg); - } - - /** - * 发送动态航班信息 - * @param fltr - */ - protected void sendFltr(SCHD.FLTR fltr) - { - //send schd - String schdMsg = JsonUtil.getString(fltr); - kafkaservice.msgSend("schd", schdMsg); - } /** * 更新Cminmsg状态为已处理 diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FltrSendScheduled.java b/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FltrSendScheduled.java new file mode 100644 index 00000000..2502668c --- /dev/null +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FltrSendScheduled.java @@ -0,0 +1,53 @@ +package com.gzzn.omms.msgexchangeapi.scheduled; + +import java.util.List; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Async; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +import com.gzzn.omms.msgexchangeapi.entity.CminmsgWapper; +import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD; +import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR; +import com.gzzn.omms.msgexchangeapi.service.IKafkaService; +import com.gzzn.omms.msgexchangeapi.service.IMsgBufferService; +import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; + +/** + * 航班动态定时发送任务,间隔多长时间发一次 + * @author zhouxiunai + * + */ +@Component +@Async +public class FltrSendScheduled { + private static Logger logger = LoggerFactory.getLogger(FltrSendScheduled.class); + + @Autowired + private IKafkaService kafkaservice; + + @Autowired + private IMsgBufferService msgBufferService; + + @Scheduled(cron = "${scheduled.flightSendCron}") + public void scheduled() + { + // + logger.info("开始发送航班动态信息"); + + try { + List fltrs = msgBufferService.takeAllFltr(); + if(fltrs.size() > 0) + { + //send schd list + String schdMsg = JsonUtil.getString(fltrs); + kafkaservice.msgSend("schd", schdMsg); + } + } catch (Exception e) { + logger.error("动态航班信息发送异常:{}",e.getMessage()); + } + }//end shceduled +} diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/IMsgBufferService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/IMsgBufferService.java new file mode 100644 index 00000000..3aeb1ee7 --- /dev/null +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/IMsgBufferService.java @@ -0,0 +1,38 @@ +package com.gzzn.omms.msgexchangeapi.service; + +import java.util.List; + +import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD; + +/** + * 消息缓存服务,多线程安全 + * @author zhouxiunai + * + */ +public interface IMsgBufferService { + /** + * 添加动态航班信息到缓存 + * @param fltr + */ + public void addFltr(SCHD.FLTR fltr); + + /** + * 添加动态消息到缓存 + * @param msg (json 格式) + */ + public void addMsg(String msg); + + + /** + * 获取缓存中所有的动态航班信息 + * @return + */ + public List takeAllFltr(); + + + /** + * 获取缓存中的所有动态消息 + * @return + */ + public List takeAllMsg(); +} diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/MsgBufferService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/MsgBufferService.java new file mode 100644 index 00000000..3ca3ec07 --- /dev/null +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/MsgBufferService.java @@ -0,0 +1,61 @@ +package com.gzzn.omms.msgexchangeapi.service; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.ConcurrentLinkedQueue; + +import org.springframework.stereotype.Service; + +import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD; +import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR; + +/** + * 要发到kafka 的消息,先在这里做缓存。 + * @author zhouxiunai + * + */ +@Service("msgbufferservice") +public class MsgBufferService implements IMsgBufferService { + private static ConcurrentLinkedQueue fltrBuffer = new ConcurrentLinkedQueue(); + private static ConcurrentLinkedQueue msgBuffer = new ConcurrentLinkedQueue(); + + @Override + public void addFltr(SCHD.FLTR fltr) { + fltrBuffer.add(fltr); + } + + @Override + public void addMsg(String msg) { + msgBuffer.add(msg); + } + + @Override + public List takeAllFltr() { + if(fltrBuffer.size() <= 0) + { + return Collections.emptyList(); + } + + // + FLTR[] fltrs = (FLTR[]) fltrBuffer.toArray(new FLTR[fltrBuffer.size()]); // TODO 转换出错 + fltrBuffer.clear();//to fix ,获取及清空过程中是否要加锁 + + return Arrays.asList(fltrs); + } + + @Override + public List takeAllMsg() { + if(msgBuffer.size() <= 0) + { + return Collections.emptyList(); + } + + // + String[] msgs = (String[]) msgBuffer.toArray(); + msgBuffer.clear();//to fix ,获取及清空过程中是否要加锁 + + return Arrays.asList(msgs); + } + +} diff --git a/src/main/resources/application-dev.yml b/src/main/resources/application-dev.yml index 9dd9ef63..10111308 100644 --- a/src/main/resources/application-dev.yml +++ b/src/main/resources/application-dev.yml @@ -65,6 +65,8 @@ scheduled: queueCapacity: 10 # 每日凌晨3:00 cminmsg 已处理数据转历史 cminmsgHisCron: 0 0 3 * * ? + # 航班发送信息定时 + flightSendCron: "*/5 * * * * ?" #计入历史航班数据的条件 hstCondition: #航班到达超过多久则成为历史航班数据(单位:秒)