package com.gzzn.omms.msgexchangeapi.scheduled; import java.util.ArrayList; import java.util.List; import java.util.Locale; 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; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; 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; import com.gzzn.omms.msgexchangeapi.utils.DateTimeUtil; import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; /** * 历史航班数据定时器 * @author fhc * */ @Component @Async public class FlightHisScheduled { @Autowired private IFlightInfoService flightInfoService; /** * 航班到达超过多久则成为历史航班数据(单位:秒) */ @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 ; /** * 历史航班数据索引 */ private final String INDEX_NAME = "flight_hts"; /** * 历史航班数据类型 */ private final String ES_TYPE = "_doc"; private final Logger logger = LoggerFactory.getLogger(FlightHisScheduled.class); /** * 每天凌晨 将前一天的 动态消息中 已为历史的 放入ES中 */ @Scheduled(cron = "${scheduled.flightInfoCron}") public void scheduled(){ logger.info("开始采集历史数据..."); //1.从redis中获取 当前的动态消息 List fltrs = flightInfoService.findAll(); List fltrIdForHistory = new ArrayList(); List fltrsHistory = new ArrayList(); //遍历 查找出已经 可以作为历史数据的动态消息 //条件为 :完成运营的航班是指已经落地超过2小时或已经起飞的 if(null!=fltrs && fltrs.size()>0){ for(FLTR fltr : fltrs){ 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; } /**===================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; } } Long todayMixTime = DateTimeUtil.getTodayStartTime(); String ACTT = fltr.getACTT(); //航班实际时间 ddMMMyyHHmm if(!StringUtils.isEmpty(ACTT) && SODT_longARRIVE_HST_TIME){ //转为历史数据 fltrsHistory.add(fltr); continue; } } } } //往 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().toString()); }else{ //新增历史数据 ElasticsearchUtil.addData(jsondata, INDEX_NAME, ES_TYPE); fltrIdForHistory.add(fltr.getFLID().toString()); } } } //删除redis中已成为历史数据的动态消息 if(null!=fltrIdForHistory && fltrIdForHistory.size()>0){ flightInfoService.batchDeleteFltr(fltrIdForHistory); } logger.info("结束采集历史数据。"); } } }