动态转历史 需求调整

This commit is contained in:
fanghongchao
2018-12-19 17:11:32 +08:00
parent 7e33b4e0a3
commit 48e4ad2053
2 changed files with 109 additions and 60 deletions
@@ -8,6 +8,8 @@ import java.util.Map;
import org.elasticsearch.index.query.BoolQueryBuilder;
import org.elasticsearch.index.query.QueryBuilders;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.scheduling.annotation.Async;
@@ -17,6 +19,7 @@ import org.springframework.util.StringUtils;
import com.alibaba.fastjson.JSONObject;
import com.gzzn.omms.msgexchangeapi.elasticsearch.ElasticsearchUtil;
import com.gzzn.omms.msgexchangeapi.entity.msg.DIVERSIONDATA;
import com.gzzn.omms.msgexchangeapi.entity.msg.MOVEMENTINDICATOR;
import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD.FLTR;
import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService;
@@ -40,6 +43,23 @@ public class FlightHisScheduled {
*/
@Value("${hstCondition.ARRIVE_HST_TIME}")
private Integer ARRIVE_HST_TIME ;
/**
* 航班取消超过多久则成为历史航班数据(单位:秒)
*/
@Value("${hstCondition.CANCEL_HST_TIME}")
private Integer CANCEL_HST_TIME ;
/**
* 备降超过多久则成为历史航班数据 FROM CTU 本应降落到成都备降到其他地方(单位:秒)
*/
@Value("${hstCondition.FDIV_HST_TIME}")
private Integer FDIV_HST_TIME ;
/**
* 计划时间超过多久则成为历史航班数据(单位:秒)
*/
@Value("${hstCondition.SODT_HST_TIME}")
private Integer SODT_HST_TIME ;
/**
* 历史航班数据索引
*/
@@ -49,89 +69,113 @@ public class FlightHisScheduled {
*/
private final String ES_TYPE = "_doc";
private final Logger logger = LoggerFactory.getLogger(FlightHisScheduled.class);
/**
* 每天凌晨 将前一天的 动态消息中 已为历史的 放入ES中
*/
@Scheduled(cron = "${scheduled.flightInfoCron}")
public void scheduled(){
System.out.println("开始采集历史数据...");
logger.info("开始采集历史数据...");
//1.从redis中获取 当前的动态消息
List<FLTR> fltrs = flightInfoService.findAll();
List<BigInteger> fltrIdForHistory = new ArrayList<BigInteger>();
List<FLTR> fltrsHistory = new ArrayList<FLTR>();
Long todayMixTime = DateTimeUtil.getTodayStartTime();
//遍历 查找出已经 可以作为历史数据的动态消息
//条件为 :完成运营的航班是指已经落地超过2小时或已经起飞的
if(null!=fltrs && fltrs.size()>0){
for(FLTR fltr : fltrs){
String ACTT = fltr.getACTT(); //航班实际时间 ddMMMyyHHmm
if(StringUtils.isEmpty(ACTT)){
String SODT = fltr.getSODT(); //航班计划时间 ddMMMyyHHmm
Long SODT_long = DateTimeUtil.toDate(SODT, "ddMMMyyHHmm",Locale.ENGLISH).getTime();
/**===================1.计划时间超时的===================**/
if((System.currentTimeMillis()-SODT_long)>SODT_HST_TIME){
//转为历史数据
fltrsHistory.add(fltr);
continue;
}
Long ACTT_longTime = DateTimeUtil.toDate(ACTT, "ddMMMyyHHmm",Locale.ENGLISH).getTime();
if(ACTT_longTime>=todayMixTime){
//计划时间为凌晨0点(包含)后的 不计入历史数据
continue;
}
//离港航班
if(fltr.getMVIN().equals(MOVEMENTINDICATOR.D)){
//实际时间小于当前时间则为已离港
if(ACTT_longTime<=System.currentTimeMillis()){
/**===================2.航班取消===================**/
String CNCL = fltr.getCNCL(); //航班取消时间 ddMMMyyHHmm
if(!StringUtils.isEmpty(CNCL)){ //航班被取消
Long CNCL_longTime = DateTimeUtil.toDate(CNCL, "ddMMMyyHHmm",Locale.ENGLISH).getTime();
if((System.currentTimeMillis()-CNCL_longTime)>CANCEL_HST_TIME){
//转为历史数据
fltrsHistory.add(fltr);
continue;
}
}
/**===================3.备降===================**/
DIVERSIONDATA FDIV = fltr.getFDIV(); //备降信息
if(null!=FDIV && FDIV.getDDES().equals("CTU") && FDIV.getDDIR().equals("FROM")){ //本应降落到成都机场 备降到其他机场
if((System.currentTimeMillis()-SODT_long)>FDIV_HST_TIME){
//转为历史数据
fltrsHistory.add(fltr);
continue;
}
}
//到达航班
if(fltr.getMVIN().equals(MOVEMENTINDICATOR.A)){
//到达实际时间超过规定时间的 则视为历史数据
if((System.currentTimeMillis()-ACTT_longTime)>ARRIVE_HST_TIME){
//转为历史数据
fltrsHistory.add(fltr);
Long todayMixTime = DateTimeUtil.getTodayStartTime();
String ACTT = fltr.getACTT(); //航班实际时间 ddMMMyyHHmm
if(!StringUtils.isEmpty(ACTT) && SODT_long<todayMixTime){ //计划时间是今天之前 且 实际时间存在
Long ACTT_longTime = DateTimeUtil.toDate(ACTT, "ddMMMyyHHmm",Locale.ENGLISH).getTime();
/**===================4.离港航班===================**/
if(fltr.getMVIN().equals(MOVEMENTINDICATOR.D)){
//实际时间且小于当前时间则为已离港
if(ACTT_longTime<=System.currentTimeMillis()){
//转为历史数据
fltrsHistory.add(fltr);
continue;
}
}
/**===================5.到达航班===================**/
if(fltr.getMVIN().equals(MOVEMENTINDICATOR.A)){
//到达实际时间超过规定时间的 则视为历史数据
if((System.currentTimeMillis()-ACTT_longTime)>ARRIVE_HST_TIME){
//转为历史数据
fltrsHistory.add(fltr);
continue;
}
}
}
}
}
//往 elasticsearch中放入历史数据 通过FLID + ACTT 判断唯一
if(!ElasticsearchUtil.isIndexExist(INDEX_NAME)){
ElasticsearchUtil.createIndex(INDEX_NAME);
}
if(null!=fltrsHistory && fltrsHistory.size()>0){
for(FLTR fltr:fltrsHistory){
JSONObject jsondata = JSONObject.parseObject(JsonUtil.getString(fltr));
//判断是否已经拥有该历史数据 因es类型原因请注意类型转换
String ACTT = fltr.getACTT();
Long FLID = fltr.getFLID().longValue();
BoolQueryBuilder boolQuery = QueryBuilders.boolQuery()
.must(QueryBuilders.termQuery("ACTT",ACTT))
.must(QueryBuilders.termQuery("FLID",FLID));
List<Map<String, Object>> listInEs = ElasticsearchUtil.searchListData(INDEX_NAME, ES_TYPE, boolQuery, null, "FLID", null, null, null);
if(null!=listInEs && listInEs.size()>0){
//获取es中的id 并更新当前历史数据
//唯一 获取第一个
Map<String, Object> map = listInEs.get(0);
String id = (String) map.get("id");
ElasticsearchUtil.updateDataById(jsondata, INDEX_NAME, ES_TYPE, id);
fltrIdForHistory.add(fltr.getFLID());
}else{
//新增历史数据
ElasticsearchUtil.addData(jsondata, INDEX_NAME, ES_TYPE);
fltrIdForHistory.add(fltr.getFLID());
//往 elasticsearch中放入历史数据 通过FLID + SODT 判断唯一 (航班ID+计划时间)
if(!ElasticsearchUtil.isIndexExist(INDEX_NAME)){
ElasticsearchUtil.createIndex(INDEX_NAME);
}
if(null!=fltrsHistory && fltrsHistory.size()>0){
for(FLTR fltr:fltrsHistory){
JSONObject jsondata = JSONObject.parseObject(JsonUtil.getString(fltr));
//判断是否已经拥有该历史数据 因es类型原因请注意类型转换
String SODT = fltr.getSODT();
Long FLID = fltr.getFLID().longValue();
BoolQueryBuilder boolQuery = QueryBuilders.boolQuery()
.must(QueryBuilders.termQuery("SODT",SODT))
.must(QueryBuilders.termQuery("FLID",FLID));
List<Map<String, Object>> listInEs = ElasticsearchUtil.searchListData(INDEX_NAME, ES_TYPE, boolQuery, null, "FLID", null, null, null);
if(null!=listInEs && listInEs.size()>0){
//获取es中的id 并更新当前历史数据
//唯一 获取第一个
Map<String, Object> map = listInEs.get(0);
String id = (String) map.get("id");
ElasticsearchUtil.updateDataById(jsondata, INDEX_NAME, ES_TYPE, id);
fltrIdForHistory.add(fltr.getFLID());
}else{
//新增历史数据
ElasticsearchUtil.addData(jsondata, INDEX_NAME, ES_TYPE);
fltrIdForHistory.add(fltr.getFLID());
}
}
}
//删除redis中已成为历史数据的动态消息
if(null!=fltrIdForHistory && fltrIdForHistory.size()>0){
flightInfoService.batchDeleteFltr(fltrIdForHistory);
}
logger.info("结束采集历史数据。");
}
//删除redis中已成为历史数据的动态消息
if(null!=fltrIdForHistory && fltrIdForHistory.size()>0){
flightInfoService.batchDeleteFltr(fltrIdForHistory);
}
System.out.println("结束采集历史数据...");
}
}
}
+9 -4
View File
@@ -54,16 +54,21 @@ msgExchange:
logstash:
host: 130.120.3.233:5000
scheduled:
#每日凌晨2点采集历史数据
#flightInfoCron: 0 14 18 * * ?
flightInfoCron: 0 0 2 * * ?
#每日凌晨3:30分点采集历史数据 0 30 3 * * ?
flightInfoCron: 0 30 3 * * ?
corePoolSize: 10
maxPoolSize: 50
queueCapacity: 10
#计入历史航班数据的条件
hstCondition:
#航班到达超过多久则成为历史航班数据(单位:秒)
ARRIVE_HST_TIME: 7200
ARRIVE_HST_TIME: 3600
#航班取消超过多久则成为历史航班数据(单位:秒)
CANCEL_HST_TIME: 3600
#备降超过多久则成为历史航班数据 FROM CTU 本应降落到成都备降到其他地方(单位:秒)
FDIV_HST_TIME: 59400
#计划时间超过多久则成为历史航班数据(单位:秒)
SODT_HST_TIME: 259200
# Elasticsearch
# 9200端口是用来让HTTP REST API来访问ElasticSearch,而9300端口是传输层监听的默认端口