From f0dbf1a3b3224981e898454b6cc90a7ceabc8393 Mon Sep 17 00:00:00 2001 From: zhouxiunai <154707516@qq.com> Date: Thu, 17 Jan 2019 20:59:43 +0800 Subject: [PATCH] =?UTF-8?q?=E9=85=8D=E5=90=88=E5=89=8D=E7=AB=AF=E5=87=8F?= =?UTF-8?q?=E5=B0=91=E6=B8=B2=E6=9F=93=E6=AC=A1=E6=95=B0=EF=BC=8C=E5=B0=86?= =?UTF-8?q?=E5=8A=A8=E6=80=81=E8=88=AA=E7=8F=AD=E4=BF=A1=E6=81=AF=E5=8A=A0?= =?UTF-8?q?=E5=85=A5=E5=88=B0=E7=BC=93=E5=AD=98=E3=80=82=E9=9A=94=E6=8C=87?= =?UTF-8?q?=E5=AE=9A=E6=97=B6=E9=97=B4xx=20=E5=90=8E=EF=BC=8C=E4=BB=A5?= =?UTF-8?q?=E6=95=B0=E7=BB=84=E7=9A=84=E5=BD=A2=E5=BC=8F=E4=B8=80=E6=AC=A1?= =?UTF-8?q?=E5=8F=91=E5=87=BA=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../msghandler/flop/FDELHandler.java | 2 +- .../msghandler/flop/base/FlopBaseHandler.java | 45 +++++--------- .../msghandler/schd/ADFTHandler.java | 5 +- .../msghandler/schd/base/SchdBaseHandler.java | 42 ++++--------- .../scheduled/FltrSendScheduled.java | 53 ++++++++++++++++ .../service/IMsgBufferService.java | 38 ++++++++++++ .../service/MsgBufferService.java | 61 +++++++++++++++++++ src/main/resources/application-dev.yml | 2 + 8 files changed, 185 insertions(+), 63 deletions(-) create mode 100644 src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FltrSendScheduled.java create mode 100644 src/main/java/com/gzzn/omms/msgexchangeapi/service/IMsgBufferService.java create mode 100644 src/main/java/com/gzzn/omms/msgexchangeapi/service/MsgBufferService.java 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: #航班到达超过多久则成为历史航班数据(单位:秒)