138 lines
4.7 KiB
Java
138 lines
4.7 KiB
Java
package com.gzzn.omms.msgexchangeapi.scheduled;
|
|||
|
|
|
||
|
|
import java.math.BigInteger;
|
||
|
|
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.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.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 ;
|
||
|
|
/**
|
||
|
|
* 历史航班数据索引
|
||
|
|
*/
|
||
|
|
private final String INDEX_NAME = "flight_hts";
|
||
|
|
/**
|
||
|
|
* 历史航班数据类型
|
||
|
|
*/
|
||
|
|
private final String ES_TYPE = "_doc";
|
||
|
|
|
||
|
|
|
||
|
|
/**
|
||
|
|
* 每天凌晨 将前一天的 动态消息中 已为历史的 放入ES中
|
||
|
|
*/
|
||
|
|
@Scheduled(cron = "${scheduled.flightInfoCron}")
|
||
|
|
public void scheduled(){
|
||
|
|
System.out.println("开始采集历史数据...");
|
||
|
|
//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)){
|
||
|
|
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()){
|
||
|
|
//转为历史数据
|
||
|
|
fltrsHistory.add(fltr);
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
//到达航班
|
||
|
|
if(fltr.getMVIN().equals(MOVEMENTINDICATOR.A)){
|
||
|
|
//到达实际时间超过规定时间的 则视为历史数据
|
||
|
|
if((System.currentTimeMillis()-ACTT_longTime)>ARRIVE_HST_TIME){
|
||
|
|
//转为历史数据
|
||
|
|
fltrsHistory.add(fltr);
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
|
||
|
|
//往 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());
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
//删除redis中已成为历史数据的动态消息
|
||
|
|
if(null!=fltrIdForHistory && fltrIdForHistory.size()>0){
|
||
|
|
flightInfoService.batchDeleteFltr(fltrIdForHistory);
|
||
|
|
}
|
||
|
|
|
||
|
|
System.out.println("结束采集历史数据...");
|
||
|
|
}
|
||
|
|
|
||
|
|
}
|