重构航班缓存发送代码
This commit is contained in:
+3
-3
@@ -15,7 +15,7 @@ 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.IFltrSendBufferService;
|
||||||
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.SpringUtil;
|
import com.gzzn.omms.msgexchangeapi.utils.SpringUtil;
|
||||||
@@ -33,7 +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;
|
protected IFltrSendBufferService msgbufferservice;
|
||||||
|
|
||||||
public FlopBaseHandler()
|
public FlopBaseHandler()
|
||||||
{
|
{
|
||||||
@@ -42,7 +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");
|
msgbufferservice = (IFltrSendBufferService) SpringUtil.getBean("msgbufferservice");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+3
-3
@@ -13,7 +13,7 @@ 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.IFltrSendBufferService;
|
||||||
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.SpringUtil;
|
import com.gzzn.omms.msgexchangeapi.utils.SpringUtil;
|
||||||
@@ -30,7 +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;
|
protected IFltrSendBufferService msgbufferservice;
|
||||||
|
|
||||||
public SchdBaseHandler()
|
public SchdBaseHandler()
|
||||||
{
|
{
|
||||||
@@ -38,7 +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");
|
msgbufferservice = (IFltrSendBufferService) SpringUtil.getBean("msgbufferservice");
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|||||||
@@ -9,11 +9,9 @@ import org.springframework.scheduling.annotation.Async;
|
|||||||
import org.springframework.scheduling.annotation.Scheduled;
|
import org.springframework.scheduling.annotation.Scheduled;
|
||||||
import org.springframework.stereotype.Component;
|
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.entity.msg.SCHD.FLTR;
|
||||||
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.IFltrSendBufferService;
|
||||||
import com.gzzn.omms.msgexchangeapi.utils.JsonUtil;
|
import com.gzzn.omms.msgexchangeapi.utils.JsonUtil;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -30,7 +28,7 @@ public class FltrSendScheduled {
|
|||||||
private IKafkaService kafkaservice;
|
private IKafkaService kafkaservice;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
private IMsgBufferService msgBufferService;
|
private IFltrSendBufferService msgBufferService;
|
||||||
|
|
||||||
@Scheduled(cron = "${scheduled.flightSendCron}")
|
@Scheduled(cron = "${scheduled.flightSendCron}")
|
||||||
public void scheduled()
|
public void scheduled()
|
||||||
|
|||||||
+2
-22
@@ -11,24 +11,19 @@ import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD;
|
|||||||
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR;
|
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 要发到kafka 的消息,先在这里做缓存。
|
* 要发到kafka 的航班动态信息,先在这里做缓存。
|
||||||
* @author zhouxiunai
|
* @author zhouxiunai
|
||||||
*
|
*
|
||||||
*/
|
*/
|
||||||
@Service("msgbufferservice")
|
@Service("msgbufferservice")
|
||||||
public class MsgBufferService implements IMsgBufferService {
|
public class FltrSendBufferService implements IFltrSendBufferService {
|
||||||
private static ConcurrentLinkedQueue<SCHD.FLTR> fltrBuffer = new ConcurrentLinkedQueue<SCHD.FLTR>();
|
private static ConcurrentLinkedQueue<SCHD.FLTR> fltrBuffer = new ConcurrentLinkedQueue<SCHD.FLTR>();
|
||||||
private static ConcurrentLinkedQueue<String> msgBuffer = new ConcurrentLinkedQueue<String>();
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void addFltr(SCHD.FLTR fltr) {
|
public void addFltr(SCHD.FLTR fltr) {
|
||||||
fltrBuffer.add(fltr);
|
fltrBuffer.add(fltr);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
|
||||||
public void addMsg(String msg) {
|
|
||||||
msgBuffer.add(msg);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public List<FLTR> takeAllFltr() {
|
public List<FLTR> takeAllFltr() {
|
||||||
@@ -43,19 +38,4 @@ public class MsgBufferService implements IMsgBufferService {
|
|||||||
|
|
||||||
return Arrays.asList(fltrs);
|
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);
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
}
|
||||||
+2
-15
@@ -5,34 +5,21 @@ import java.util.List;
|
|||||||
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD;
|
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 消息缓存服务,多线程安全
|
* 航班信息发送缓存服务,多线程安全
|
||||||
* @author zhouxiunai
|
* @author zhouxiunai
|
||||||
*
|
*
|
||||||
*/
|
*/
|
||||||
public interface IMsgBufferService {
|
public interface IFltrSendBufferService {
|
||||||
/**
|
/**
|
||||||
* 添加动态航班信息到缓存
|
* 添加动态航班信息到缓存
|
||||||
* @param fltr
|
* @param fltr
|
||||||
*/
|
*/
|
||||||
public void addFltr(SCHD.FLTR fltr);
|
public void addFltr(SCHD.FLTR fltr);
|
||||||
|
|
||||||
/**
|
|
||||||
* 添加动态消息到缓存
|
|
||||||
* @param msg (json 格式)
|
|
||||||
*/
|
|
||||||
public void addMsg(String msg);
|
|
||||||
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 获取缓存中所有的动态航班信息
|
* 获取缓存中所有的动态航班信息
|
||||||
* @return
|
* @return
|
||||||
*/
|
*/
|
||||||
public List<SCHD.FLTR> takeAllFltr();
|
public List<SCHD.FLTR> takeAllFltr();
|
||||||
|
|
||||||
|
|
||||||
/**
|
|
||||||
* 获取缓存中的所有动态消息
|
|
||||||
* @return
|
|
||||||
*/
|
|
||||||
public List<String> takeAllMsg();
|
|
||||||
}
|
}
|
||||||
Reference in New Issue
Block a user