配合前端减少渲染次数,将动态航班信息加入到缓存。隔指定时间xx 后,以数组的形式一次发出。
This commit is contained in:
@@ -52,7 +52,7 @@ public class FDELHandler extends FlopBaseHandler {
|
|||||||
}//end for
|
}//end for
|
||||||
|
|
||||||
//推送变更后的主航班信息
|
//推送变更后的主航班信息
|
||||||
this.sendFltr(fltrMafl);
|
this.addToBuffer(fltrMafl);
|
||||||
}
|
}
|
||||||
|
|
||||||
return HandlerResult.success(fltr);
|
return HandlerResult.success(fltr);
|
||||||
|
|||||||
+15
-30
@@ -15,9 +15,9 @@ import com.gzzn.omms.msgexchangeapi.msghandler.IBaseHandler;
|
|||||||
import com.gzzn.omms.msgexchangeapi.redis.RedisService;
|
import com.gzzn.omms.msgexchangeapi.redis.RedisService;
|
||||||
import com.gzzn.omms.msgexchangeapi.service.ICminmsgService;
|
import com.gzzn.omms.msgexchangeapi.service.ICminmsgService;
|
||||||
import com.gzzn.omms.msgexchangeapi.service.IKafkaService;
|
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.exchange.IExchangeService;
|
||||||
import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService;
|
import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService;
|
||||||
import com.gzzn.omms.msgexchangeapi.utils.JsonUtil;
|
|
||||||
import com.gzzn.omms.msgexchangeapi.utils.SpringUtil;
|
import com.gzzn.omms.msgexchangeapi.utils.SpringUtil;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -33,6 +33,7 @@ public class FlopBaseHandler implements IBaseHandler {
|
|||||||
protected IKafkaService kafkaservice;
|
protected IKafkaService kafkaservice;
|
||||||
protected RedisService redisService;
|
protected RedisService redisService;
|
||||||
protected ICminmsgService cminmsgService;
|
protected ICminmsgService cminmsgService;
|
||||||
|
protected IMsgBufferService msgbufferservice;
|
||||||
|
|
||||||
public FlopBaseHandler()
|
public FlopBaseHandler()
|
||||||
{
|
{
|
||||||
@@ -41,6 +42,7 @@ public class FlopBaseHandler implements IBaseHandler {
|
|||||||
kafkaservice = (IKafkaService) SpringUtil.getBean("kafkaservice");
|
kafkaservice = (IKafkaService) SpringUtil.getBean("kafkaservice");
|
||||||
redisService = (RedisService) SpringUtil.getBean("redisService");
|
redisService = (RedisService) SpringUtil.getBean("redisService");
|
||||||
cminmsgService = (ICminmsgService) SpringUtil.getBean("cminmsgService");
|
cminmsgService = (ICminmsgService) SpringUtil.getBean("cminmsgService");
|
||||||
|
msgbufferservice = (IMsgBufferService) SpringUtil.getBean("msgbufferservice");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -60,7 +62,8 @@ public class FlopBaseHandler implements IBaseHandler {
|
|||||||
flightInfoService.saveFltr(fltrResult);
|
flightInfoService.saveFltr(fltrResult);
|
||||||
|
|
||||||
//send and update
|
//send and update
|
||||||
this.sendFltrAndMsg(fltrResult, cminmsg);
|
this.addToBuffer(fltrResult);
|
||||||
|
this.sendMsg(cminmsg);
|
||||||
this.updateCminmsgToProcessed(cminmsg);
|
this.updateCminmsgToProcessed(cminmsg);
|
||||||
|
|
||||||
return HandlerResult.success(fltrResult);
|
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 消息
|
* @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
|
* @param cminmsg
|
||||||
|
|||||||
@@ -34,7 +34,8 @@ public class ADFTHandler extends SchdBaseHandler{
|
|||||||
flightInfoService.saveFltr(fltr);
|
flightInfoService.saveFltr(fltr);
|
||||||
|
|
||||||
|
|
||||||
this.sendFltrAndMsg(fltr, cminmsg);
|
this.addToBuffer(fltr);
|
||||||
|
this.sendMsg(cminmsg);
|
||||||
this.updateCminmsgToProcessed(cminmsg);
|
this.updateCminmsgToProcessed(cminmsg);
|
||||||
|
|
||||||
|
|
||||||
@@ -48,7 +49,7 @@ public class ADFTHandler extends SchdBaseHandler{
|
|||||||
fltrMafl.getMafl().add(mafldata);
|
fltrMafl.getMafl().add(mafldata);
|
||||||
|
|
||||||
//推送变更后的主航班信息
|
//推送变更后的主航班信息
|
||||||
this.sendFltr(fltrMafl);
|
this.addToBuffer(fltrMafl);
|
||||||
}
|
}
|
||||||
|
|
||||||
return HandlerResult.success();
|
return HandlerResult.success();
|
||||||
|
|||||||
+12
-30
@@ -5,7 +5,6 @@ import java.util.Arrays;
|
|||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
import com.gzzn.omms.msgexchangeapi.dto.ResponseDto;
|
|
||||||
import com.gzzn.omms.msgexchangeapi.entity.Cminmsg;
|
import com.gzzn.omms.msgexchangeapi.entity.Cminmsg;
|
||||||
import com.gzzn.omms.msgexchangeapi.entity.CminmsgWapper;
|
import com.gzzn.omms.msgexchangeapi.entity.CminmsgWapper;
|
||||||
import com.gzzn.omms.msgexchangeapi.entity.msg.MSG;
|
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.msghandler.IBaseHandler;
|
||||||
import com.gzzn.omms.msgexchangeapi.service.ICminmsgService;
|
import com.gzzn.omms.msgexchangeapi.service.ICminmsgService;
|
||||||
import com.gzzn.omms.msgexchangeapi.service.IKafkaService;
|
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.exchange.IExchangeService;
|
||||||
import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService;
|
import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService;
|
||||||
import com.gzzn.omms.msgexchangeapi.utils.JsonUtil;
|
|
||||||
import com.gzzn.omms.msgexchangeapi.utils.SpringUtil;
|
import com.gzzn.omms.msgexchangeapi.utils.SpringUtil;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -31,6 +30,7 @@ public class SchdBaseHandler implements IBaseHandler{
|
|||||||
protected IFlightInfoService flightInfoService;
|
protected IFlightInfoService flightInfoService;
|
||||||
protected IKafkaService kafkaservice;
|
protected IKafkaService kafkaservice;
|
||||||
protected ICminmsgService cminmsgService;
|
protected ICminmsgService cminmsgService;
|
||||||
|
protected IMsgBufferService msgbufferservice;
|
||||||
|
|
||||||
public SchdBaseHandler()
|
public SchdBaseHandler()
|
||||||
{
|
{
|
||||||
@@ -38,6 +38,7 @@ public class SchdBaseHandler implements IBaseHandler{
|
|||||||
flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService");
|
flightInfoService = (IFlightInfoService) SpringUtil.getBean("flightInfoService");
|
||||||
kafkaservice = (IKafkaService) SpringUtil.getBean("kafkaservice");
|
kafkaservice = (IKafkaService) SpringUtil.getBean("kafkaservice");
|
||||||
cminmsgService = (ICminmsgService) SpringUtil.getBean("cminmsgService");
|
cminmsgService = (ICminmsgService) SpringUtil.getBean("cminmsgService");
|
||||||
|
msgbufferservice = (IMsgBufferService) SpringUtil.getBean("msgbufferservice");
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -45,6 +46,15 @@ public class SchdBaseHandler implements IBaseHandler{
|
|||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 添加航班动态信息到缓冲区
|
||||||
|
* @param fltr
|
||||||
|
*/
|
||||||
|
protected void addToBuffer(SCHD.FLTR fltr)
|
||||||
|
{
|
||||||
|
msgbufferservice.addFltr(fltr);//缓存区
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 发送日计划到kafka队列
|
* 发送日计划到kafka队列
|
||||||
* @param cminmsg
|
* @param cminmsg
|
||||||
@@ -80,34 +90,6 @@ public class SchdBaseHandler implements IBaseHandler{
|
|||||||
String msg = exchangeService.msgToJson(cminmsg.getMsg());
|
String msg = exchangeService.msgToJson(cminmsg.getMsg());
|
||||||
kafkaservice.msgSend("msg", msg);
|
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状态为已处理
|
* 更新Cminmsg状态为已处理
|
||||||
|
|||||||
@@ -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<FLTR> 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
|
||||||
|
}
|
||||||
@@ -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<SCHD.FLTR> takeAllFltr();
|
||||||
|
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 获取缓存中的所有动态消息
|
||||||
|
* @return
|
||||||
|
*/
|
||||||
|
public List<String> takeAllMsg();
|
||||||
|
}
|
||||||
@@ -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<SCHD.FLTR> fltrBuffer = new ConcurrentLinkedQueue<SCHD.FLTR>();
|
||||||
|
private static ConcurrentLinkedQueue<String> msgBuffer = new ConcurrentLinkedQueue<String>();
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void addFltr(SCHD.FLTR fltr) {
|
||||||
|
fltrBuffer.add(fltr);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void addMsg(String msg) {
|
||||||
|
msgBuffer.add(msg);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public List<FLTR> 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<String> takeAllMsg() {
|
||||||
|
if(msgBuffer.size() <= 0)
|
||||||
|
{
|
||||||
|
return Collections.emptyList();
|
||||||
|
}
|
||||||
|
|
||||||
|
//
|
||||||
|
String[] msgs = (String[]) msgBuffer.toArray();
|
||||||
|
msgBuffer.clear();//to fix ,获取及清空过程中是否要加锁
|
||||||
|
|
||||||
|
return Arrays.asList(msgs);
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
@@ -65,6 +65,8 @@ scheduled:
|
|||||||
queueCapacity: 10
|
queueCapacity: 10
|
||||||
# 每日凌晨3:00 cminmsg 已处理数据转历史
|
# 每日凌晨3:00 cminmsg 已处理数据转历史
|
||||||
cminmsgHisCron: 0 0 3 * * ?
|
cminmsgHisCron: 0 0 3 * * ?
|
||||||
|
# 航班发送信息定时
|
||||||
|
flightSendCron: "*/5 * * * * ?"
|
||||||
#计入历史航班数据的条件
|
#计入历史航班数据的条件
|
||||||
hstCondition:
|
hstCondition:
|
||||||
#航班到达超过多久则成为历史航班数据(单位:秒)
|
#航班到达超过多久则成为历史航班数据(单位:秒)
|
||||||
|
|||||||
Reference in New Issue
Block a user