去除不必要的redis状态及消息偏移获取。默认获取所有未处理消息。 从id为0开始。

This commit is contained in:
zhouxiunai
2019-01-08 17:13:09 +08:00
parent 90d8751740
commit d425350e39
3 changed files with 1 additions and 24 deletions
@@ -8,5 +8,4 @@ package com.gzzn.omms.msgexchangeapi.redis;
public class RedisKeyConstant {
public static String KEY_ISWAITSCHD="isWaitSchd";
public static String KEY_WAITSCHDDATE="waitSchdDate";
public static String KEY_LASTBEGINID="lastBeginId";
}
@@ -11,7 +11,6 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;
import com.gzzn.omms.msgexchangeapi.redis.RedisKeyConstant;
import com.gzzn.omms.msgexchangeapi.redis.RedisService;
@Component
@@ -33,10 +32,6 @@ public class AppRunner implements CommandLineRunner {
public void run(String... args) throws Exception {
logger.info("app runner start");
//
redisService.set(RedisKeyConstant.KEY_LASTBEGINID,0L);
//启动消息采集线程
ScheduledExecutorService scheduledExecutorService = Executors.newScheduledThreadPool(1);
@@ -4,12 +4,9 @@ import java.util.List;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.util.CollectionUtils;
import com.gzzn.omms.msgexchangeapi.entity.Cminmsg;
import com.gzzn.omms.msgexchangeapi.msghandler.MsgHandlerDispatcher;
import com.gzzn.omms.msgexchangeapi.redis.RedisKeyConstant;
import com.gzzn.omms.msgexchangeapi.redis.RedisService;
import com.gzzn.omms.msgexchangeapi.service.ICminmsgService;
import com.gzzn.omms.msgexchangeapi.utils.SpringUtil;
@@ -22,11 +19,9 @@ public class MsgExchangeRunner implements Runnable{
private static Logger logger = LoggerFactory.getLogger(MsgExchangeRunner.class);
private MsgHandlerDispatcher msgHandlerDispatcher;
private RedisService redisService;
private ICminmsgService cminmsgService;
public MsgExchangeRunner() {
redisService = (RedisService) SpringUtil.getBean("redisService");
cminmsgService = (ICminmsgService) SpringUtil.getBean("cminmsgService");
msgHandlerDispatcher = (MsgHandlerDispatcher) SpringUtil.getBean("msgHandlerDispatcher");
}
@@ -36,15 +31,8 @@ public class MsgExchangeRunner implements Runnable{
public void run() {
logger.info("进入获取动态航班消息...");
//获取上次的beginId
Integer beginIdInt = (Integer)redisService.get(RedisKeyConstant.KEY_LASTBEGINID);
Long beginId = beginIdInt.longValue();
logger.info("从上次获取的消息位置{}开始",beginId);
//获取航班动态消息
List<Cminmsg> lsCminmsgs = cminmsgService.getNewMsgsAfterId(beginId);
List<Cminmsg> lsCminmsgs = cminmsgService.getNewMsgsAfterId(0L);
if(null == lsCminmsgs || lsCminmsgs.size() <= 0)
{
logger.info("本周期没有获取到最新消息,直接退出本次循环");
@@ -59,11 +47,6 @@ public class MsgExchangeRunner implements Runnable{
msgHandlerDispatcher.dispatch(cminmsg);
}//end for
if(!CollectionUtils.isEmpty(lsCminmsgs))
{
beginId = lsCminmsgs.get(lsCminmsgs.size()-1).getCminmsgsId();//偏移到最新位置,等待最新的消息
redisService.set(RedisKeyConstant.KEY_LASTBEGINID,beginId);
}
logger.info("本次定时获取新消息完成...");
}//end run
}