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 3b8329f5..223d330a 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,7 +15,7 @@ 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.IFltrSendBufferService; import com.gzzn.omms.msgexchangeapi.service.exchange.IExchangeService; import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; @@ -33,7 +33,7 @@ public class FlopBaseHandler implements IBaseHandler { protected IKafkaService kafkaservice; protected RedisService redisService; protected ICminmsgService cminmsgService; - protected IMsgBufferService msgbufferservice; + protected IFltrSendBufferService msgbufferservice; public FlopBaseHandler() { @@ -42,7 +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"); + msgbufferservice = (IFltrSendBufferService) SpringUtil.getBean("msgbufferservice"); } 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 c8c60749..d39f68b4 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 @@ -13,7 +13,7 @@ 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.IFltrSendBufferService; import com.gzzn.omms.msgexchangeapi.service.exchange.IExchangeService; import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; import com.gzzn.omms.msgexchangeapi.utils.SpringUtil; @@ -30,7 +30,7 @@ public class SchdBaseHandler implements IBaseHandler{ protected IFlightInfoService flightInfoService; protected IKafkaService kafkaservice; protected ICminmsgService cminmsgService; - protected IMsgBufferService msgbufferservice; + protected IFltrSendBufferService msgbufferservice; public SchdBaseHandler() { @@ -38,7 +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"); + msgbufferservice = (IFltrSendBufferService) SpringUtil.getBean("msgbufferservice"); } @Override diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FltrSendScheduled.java b/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FltrSendScheduled.java index 2502668c..d8f230bc 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FltrSendScheduled.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FltrSendScheduled.java @@ -9,11 +9,9 @@ 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.service.IFltrSendBufferService; import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; /** @@ -30,7 +28,7 @@ public class FltrSendScheduled { private IKafkaService kafkaservice; @Autowired - private IMsgBufferService msgBufferService; + private IFltrSendBufferService msgBufferService; @Scheduled(cron = "${scheduled.flightSendCron}") public void scheduled() diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/MsgBufferService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/FltrSendBufferService.java similarity index 59% rename from src/main/java/com/gzzn/omms/msgexchangeapi/service/MsgBufferService.java rename to src/main/java/com/gzzn/omms/msgexchangeapi/service/FltrSendBufferService.java index 3ca3ec07..d3f94c8c 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/MsgBufferService.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/FltrSendBufferService.java @@ -11,24 +11,19 @@ import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD; import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR; /** - * 要发到kafka 的消息,先在这里做缓存。 + * 要发到kafka 的航班动态信息,先在这里做缓存。 * @author zhouxiunai * */ @Service("msgbufferservice") -public class MsgBufferService implements IMsgBufferService { +public class FltrSendBufferService implements IFltrSendBufferService { 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() { @@ -43,19 +38,4 @@ public class MsgBufferService implements IMsgBufferService { 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/java/com/gzzn/omms/msgexchangeapi/service/IMsgBufferService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/IFltrSendBufferService.java similarity index 54% rename from src/main/java/com/gzzn/omms/msgexchangeapi/service/IMsgBufferService.java rename to src/main/java/com/gzzn/omms/msgexchangeapi/service/IFltrSendBufferService.java index 3aeb1ee7..6073c382 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/service/IMsgBufferService.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/IFltrSendBufferService.java @@ -5,34 +5,21 @@ import java.util.List; import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD; /** - * 消息缓存服务,多线程安全 + * 航班信息发送缓存服务,多线程安全 * @author zhouxiunai * */ -public interface IMsgBufferService { +public interface IFltrSendBufferService { /** * 添加动态航班信息到缓存 * @param fltr */ public void addFltr(SCHD.FLTR fltr); - /** - * 添加动态消息到缓存 - * @param msg (json 格式) - */ - public void addMsg(String msg); - /** * 获取缓存中所有的动态航班信息 * @return */ public List takeAllFltr(); - - - /** - * 获取缓存中的所有动态消息 - * @return - */ - public List takeAllMsg(); }