diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FlightHisScheduled.java b/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FlightHisScheduled.java index 32909abc..dd059039 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FlightHisScheduled.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/scheduled/FlightHisScheduled.java @@ -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 fltrs = flightInfoService.findAll(); List fltrIdForHistory = new ArrayList(); List fltrsHistory = new ArrayList(); - 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_longARRIVE_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> listInEs = ElasticsearchUtil.searchListData(INDEX_NAME, ES_TYPE, boolQuery, null, "FLID", null, null, null); - if(null!=listInEs && listInEs.size()>0){ - //获取es中的id 并更新当前历史数据 - //唯一 获取第一个 - Map 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> listInEs = ElasticsearchUtil.searchListData(INDEX_NAME, ES_TYPE, boolQuery, null, "FLID", null, null, null); + if(null!=listInEs && listInEs.size()>0){ + //获取es中的id 并更新当前历史数据 + //唯一 获取第一个 + Map 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("结束采集历史数据..."); - } - + } } diff --git a/src/main/resources/application-dev.yml b/src/main/resources/application-dev.yml index f57457bf..9deb1a76 100644 --- a/src/main/resources/application-dev.yml +++ b/src/main/resources/application-dev.yml @@ -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端口是传输层监听的默认端口