diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/controller/FlightsController.java b/src/main/java/com/gzzn/omms/msgexchangeapi/controller/FlightsController.java index bb67cffd..70e94b06 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/controller/FlightsController.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/controller/FlightsController.java @@ -1,15 +1,28 @@ package com.gzzn.omms.msgexchangeapi.controller; +import java.util.ArrayList; import java.util.List; +import java.util.Map; import java.util.stream.Collectors; +import org.elasticsearch.index.query.MatchAllQueryBuilder; +import org.elasticsearch.index.query.QueryBuilders; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; +import com.alibaba.fastjson.JSONObject; +import com.fasterxml.jackson.databind.ObjectMapper; import com.gzzn.omms.msgexchangeapi.dto.ResponseDto; +import com.gzzn.omms.msgexchangeapi.elasticsearch.ElasticsearchUtil; +import com.gzzn.omms.msgexchangeapi.entity.msg.ABNDATA; +import com.gzzn.omms.msgexchangeapi.entity.msg.DIVERSIONDATA; +import com.gzzn.omms.msgexchangeapi.entity.msg.OPTDELAYDATA; +import com.gzzn.omms.msgexchangeapi.entity.msg.RETURNDATA; import com.gzzn.omms.msgexchangeapi.entity.msg.SCHD; import com.gzzn.omms.msgexchangeapi.service.flightInfo.IFlightInfoService; +import com.gzzn.omms.msgexchangeapi.utils.JsonUtil; @RestController public class FlightsController { @@ -32,4 +45,86 @@ public class FlightsController { return ResponseDto.success(fltrsRt); } + + /** + * 迁移航班历史信息 + * @return + */ + @GetMapping("/flights/migrate") + public ResponseDto migrateFlightHst() + { + //get old fltrs + MatchAllQueryBuilder matchAllQueryBuilder = QueryBuilders.matchAllQuery(); + List> listInEs = ElasticsearchUtil + .searchListDataWithScoll( + "flight_hts", + "_doc", + matchAllQueryBuilder, + null, + null, + null); + //trans to new fltrs + List lsFltrs = transToNewFltrs(listInEs); + + //fltrs save to new index + for (SCHD.FLTR fltr : lsFltrs) { + JSONObject jsondata = JSONObject.parseObject(JsonUtil.getString(fltr)); + ElasticsearchUtil.addData(jsondata, "flight_hts2", "_doc"); + } + + return ResponseDto.success(); + }//end function + + + + private List transToNewFltrs(List> listInEs) + { + List lsFltrs = new ArrayList(); + for (Map map : listInEs) { + ObjectMapper objMappper = new ObjectMapper(); + SCHD.FLTR fltr = objMappper.convertValue(map, SCHD.FLTR.class); + + //add BODY to Spec key + for (int i=0;i < fltr.getAbn().size();i++) { + ABNDATA abn = fltr.getAbn().get(i); + + List> lsAbn = (List>)map.get("ABN"); + Map abnMap = lsAbn.get(i); + + switch (abn.getType()) { + case DLY: + abn.setDly(mapToDly((List>) abnMap.get("BODY"))); + break; + case CAN: + abn.setCan(objMappper.convertValue(abnMap.get("BODY"), String.class)); + break; + case RTN: + abn.setRtn(objMappper.convertValue(abnMap.get("BODY"), RETURNDATA.class)); + break; + case ALT: + abn.setAlt(objMappper.convertValue(abnMap.get("BODY"), DIVERSIONDATA .class)); + break; + } + } //end for + + lsFltrs.add(fltr); + }//end for + + return lsFltrs; + }//end function + + + private List mapToDly(List> ls) + { + List rtLs = new ArrayList<>(ls.size()); + + for(Map map : ls) + { + ObjectMapper objMappper = new ObjectMapper(); + OPTDELAYDATA optdelaydata = objMappper.convertValue(map, OPTDELAYDATA.class); + rtLs.add(optdelaydata); + } + + return rtLs; + } } diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/elasticsearch/ElasticsearchUtil.java b/src/main/java/com/gzzn/omms/msgexchangeapi/elasticsearch/ElasticsearchUtil.java index 40528824..c90d4349 100644 --- a/src/main/java/com/gzzn/omms/msgexchangeapi/elasticsearch/ElasticsearchUtil.java +++ b/src/main/java/com/gzzn/omms/msgexchangeapi/elasticsearch/ElasticsearchUtil.java @@ -20,6 +20,7 @@ import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.action.update.UpdateRequest; import org.elasticsearch.client.transport.TransportClient; import org.elasticsearch.common.text.Text; +import org.elasticsearch.common.unit.TimeValue; import org.elasticsearch.index.query.QueryBuilder; import org.elasticsearch.search.SearchHit; import org.elasticsearch.search.fetch.subphase.highlight.HighlightBuilder; @@ -125,7 +126,8 @@ public class ElasticsearchUtil { LOGGER.info("addData response status:{},id:{}", response.status().getStatus(), response.getId()); return response.getId(); } - + + /** * 数据添加 * @@ -260,8 +262,82 @@ public class ElasticsearchUtil { return null; } + + + + /** + * 使用分词查询 + * + * @param index 索引名称 + * @param type 类型名称,可传入多个type逗号分隔 + * @param query 查询条件 + * @param fields 需要显示的字段,逗号分隔(缺省为全部字段) + * @param sortField 排序字段 + * @param sortOrder + * @param highlightField 高亮字段 + * @return + */ + public static List> searchListDataWithScoll( + String index, String type, QueryBuilder query, + String fields, String sortField, SortOrder sortOrder) { + SearchRequestBuilder searchRequestBuilder = client.prepareSearch(index); + + if (!StringUtils.isEmpty(type)) { + searchRequestBuilder.setTypes(type.split(",")); + } + + if (!StringUtils.isEmpty(fields)) { + searchRequestBuilder.setFetchSource(fields.split(","), null); + } + + if(null==sortOrder){ + sortOrder = SortOrder.DESC; + } + + if (!StringUtils.isEmpty(sortField)) { + searchRequestBuilder.addSort(sortField, sortOrder); + } + + int size = 1000; + searchRequestBuilder.setQuery(query); + searchRequestBuilder.setFetchSource(true); + searchRequestBuilder.setSize(size).setScroll(new TimeValue(2000)); + + + //打印的内容 可以在 Elasticsearch head 和 Kibana 上执行查询 + LOGGER.info("\n{}", searchRequestBuilder); + + List> sourceList = new ArrayList>(); + SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); + sourceList.addAll(getSearchResponse(searchResponse)); + + //根据total 滚动搜索 + long page = searchResponse.getHits().getTotalHits() / size; + for (int i = 0; i < page; i++) { + //再次发送请求,并使用上次搜索结果的ScrollId + searchResponse = client.prepareSearchScroll(searchResponse.getScrollId()) + .setScroll(new TimeValue(20000)).execute() + .actionGet(); + + sourceList.addAll(getSearchResponse(searchResponse)); + } + + return sourceList; + } + private static List> getSearchResponse(SearchResponse searchResponse) + { + List> sourceList = new ArrayList>(); + for (SearchHit searchHit : searchResponse.getHits().getHits()) { + searchHit.getSourceAsMap().put("id", searchHit.getId()); + sourceList.add(searchHit.getSourceAsMap()); + } + + return sourceList; + } + + /** * 高亮结果集 特殊处理 * diff --git a/src/test/java/com/gzzn/omms/msgexchangeapi/controller/FlightsControllerTest.java b/src/test/java/com/gzzn/omms/msgexchangeapi/controller/FlightsControllerTest.java index 690c5a8e..f2d4be45 100644 --- a/src/test/java/com/gzzn/omms/msgexchangeapi/controller/FlightsControllerTest.java +++ b/src/test/java/com/gzzn/omms/msgexchangeapi/controller/FlightsControllerTest.java @@ -36,4 +36,20 @@ public class FlightsControllerTest { .andExpect(jsonPath("$.is_success").value(true)) .andDo(print()); } + + + /** + * 测试获取所有航班信息接口 + * @throws Exception + */ + @Test + public void testMigrateFilghts() throws Exception + { + mvc. + perform(get("/flights/migrate")) + .andExpect(status().isOk()) + .andExpect(content().contentType(MediaType.APPLICATION_JSON_UTF8)) + .andExpect(jsonPath("$.is_success").value(true)) + .andDo(print()); + } }