skjcmanager/skjcmanager-service-api/skjcmanager-nky-api/src/main/java/cn/gistack/nky/fegin/IHsybClient.java
New file @@ -0,0 +1,26 @@ package cn.gistack.nky.fegin; import org.springframework.cloud.openfeign.FeignClient; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestParam; import java.util.List; @FeignClient( value = "blade-nky", fallback = IHsybClientFallback.class ) public interface IHsybClient { String API_PREFIX = "/client"; String UPDATE_CJYSJ_BY_TM= API_PREFIX + "/updateCjysjByTm"; String GET_ALL_RES= API_PREFIX + "/getAllRes"; @PostMapping(UPDATE_CJYSJ_BY_TM) String updateCjysjByTm(@RequestBody List<String> resIds); } skjcmanager/skjcmanager-service-api/skjcmanager-nky-api/src/main/java/cn/gistack/nky/fegin/IHsybClientFallback.java
New file @@ -0,0 +1,18 @@ package cn.gistack.nky.fegin; import cn.gistack.nky.entity.AlarmGet; import cn.gistack.nky.vo.AlarmGetVO; import cn.gistack.nky.vo.PageVO; import org.springframework.stereotype.Component; import java.util.List; @Component public class IHsybClientFallback implements IHsybClient{ @Override public String updateCjysjByTm(List<String> resIds) { return ""; } } skjcmanager/skjcmanager-service-api/skjcmanager-nky-api/src/main/java/cn/gistack/nky/resultpojo/HsybGetFuturePo.java
@@ -7,7 +7,7 @@ @Data public class HsybGetFuturePo { List<List<String>> data; Object data; String respCode; skjcmanager/skjcmanager-service-api/skjcmanager-nky-api/src/main/java/cn/gistack/nky/resultpojo/NewSwFuturePo.java
@@ -4,8 +4,12 @@ @Data public class NewSwFuturePo { private String time; //水位 private String sw; //预测时间 private String time; //水库id private String resId; } skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/sjztmd/feign/IAttResBaseClient.java
New file @@ -0,0 +1,22 @@ package cn.gistack.sm.sjztmd.feign; import org.springblade.core.tool.api.R; import org.springframework.cloud.openfeign.FeignClient; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import java.util.List; @FeignClient( value = "blade-sm", fallback = IAttResBaseClientFallback.class ) public interface IAttResBaseClient { String API_PREFIX = "/client"; String GET_ALL_RES = API_PREFIX + "/getAllRes"; @GetMapping(GET_ALL_RES) List<String> getAllRes(); } skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/sjztmd/feign/IAttResBaseClientFallback.java
New file @@ -0,0 +1,24 @@ package cn.gistack.sm.sjztmd.feign; import cn.gistack.sm.sjztmd.vo.PersonVO; import org.springblade.core.tool.api.R; import org.springframework.stereotype.Component; import java.util.List; import java.util.Set; /** * @ClassName Feign失败配置 * @Description TODO * @Author aix * @Date 2023/4/23 19:40 * @Version 1.0 */ @Component public class IAttResBaseClientFallback implements IAttResBaseClient{ @Override public List<String> getAllRes() { return null; } } skjcmanager/skjcmanager-service/skjcmanager-nky/src/main/java/cn/gistack/nky/feign/HsybClientImpl.java
New file @@ -0,0 +1,33 @@ package cn.gistack.nky.feign; import cn.gistack.nky.fegin.IHsybClient; import cn.gistack.nky.service.IHsybService; import lombok.AllArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springblade.core.tenant.annotation.NonDS; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RestController; import springfox.documentation.annotations.ApiIgnore; import java.util.List; @NonDS @ApiIgnore @RestController @AllArgsConstructor @Slf4j public class HsybClientImpl implements IHsybClient { private IHsybService hsybService; @Override @PostMapping(UPDATE_CJYSJ_BY_TM) public String updateCjysjByTm(List<String> resIds) { return hsybService.updateCjysjByTm(resIds); } } skjcmanager/skjcmanager-service/skjcmanager-nky/src/main/java/cn/gistack/nky/service/IHsybService.java
@@ -9,4 +9,16 @@ List<List<String>> getFuture(String resCd); List<NewSwFuturePo> getNewFuture(String resId); /** * 同步降雨数据 * @param resIds 水库id * @return */ String updateCjysjByTm(List<String> resIds); } skjcmanager/skjcmanager-service/skjcmanager-nky/src/main/java/cn/gistack/nky/service/impl/HstPredictServiceImpl.java
@@ -37,12 +37,12 @@ @Override public void getData(String type) { List<DataResChildrenPo> cdList = getCdList(type); List<DataResChildrenPo> collect = cdList.stream().filter( // item -> item.getRes_cd().equals("42011640018") || item.getRes_cd().equals("42128140006") item -> item.getRes_cd().equals("42011640018") || item.getRes_cd().equals("42130350046")|| item.getRes_cd().equals("42092250024")|| item.getRes_cd().equals("42022250039")|| item.getRes_cd().equals("42132150292") ).collect(Collectors.toList()); collect.forEach(dataResChildrenPo -> { // List<DataResChildrenPo> collect = cdList.stream().filter( //// item -> item.getRes_cd().equals("42011640018") || item.getRes_cd().equals("42128140006") // item -> item.getRes_cd().equals("42011640018") || item.getRes_cd().equals("42130350046")|| item.getRes_cd().equals("42092250024")|| item.getRes_cd().equals("42022250039")|| item.getRes_cd().equals("42132150292") // // ).collect(Collectors.toList()); cdList.forEach(dataResChildrenPo -> { //循环请求预测结果 QueryReqPo queryReqPo = new QueryReqPo(dataResChildrenPo, type); //请求预测结果 @@ -60,13 +60,13 @@ public void predictData(String type) { List<DataResChildrenPo> cdList = getCdList(type); //暂时筛选出数据 List<DataResChildrenPo> collect = cdList.stream().filter( item -> item.getRes_cd().equals("42011640018") || item.getRes_cd().equals("42130350046")|| item.getRes_cd().equals("42092250024")|| item.getRes_cd().equals("42022250039")|| item.getRes_cd().equals("42132150292") ).collect(Collectors.toList()); // List<DataResChildrenPo> collect = cdList.stream().filter( // item -> item.getRes_cd().equals("42011640018") || item.getRes_cd().equals("42130350046")|| item.getRes_cd().equals("42092250024")|| item.getRes_cd().equals("42022250039")|| item.getRes_cd().equals("42132150292") // ).collect(Collectors.toList()); // predict(collect,type); newPredict(collect,type); newPredict(cdList,type); } skjcmanager/skjcmanager-service/skjcmanager-nky/src/main/java/cn/gistack/nky/service/impl/HsybServiceImpl.java
@@ -1,19 +1,25 @@ package cn.gistack.nky.service.impl; import cn.gistack.common.utils.CommonUtil; import cn.gistack.common.utils.HttpClientUtils; import cn.gistack.common.utils.SpringContextUtil; import cn.gistack.nky.resultpojo.*; import cn.gistack.nky.service.IHsybService; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.parser.Feature; import lombok.AllArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springblade.core.tool.utils.DateUtil; import org.springblade.core.tool.utils.StringUtil; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.*; import org.springframework.stereotype.Service; import org.springframework.util.MultiValueMap; import org.springframework.web.client.RestTemplate; import java.time.LocalDate; import java.time.format.DateTimeFormatter; import java.util.*; import java.util.stream.Collectors; @@ -22,12 +28,22 @@ @Slf4j public class HsybServiceImpl implements IHsybService { private RestTemplate restTemplate; private static String PREFIX = "/hsybApi"; private static String ONLINE = "http://10.42.7.148:50001"; private static String LOCAL = "https://sk.hubeishuiyi.cn"; private static String GET_FUTURE = "/hsybApi/api/fh-admin/skkr/getFuture"; private static String GET_NEW_FUTURE = "/hsybApi/api/fh-admin/skkr/getSkFutureSw"; private static String GET_FUTURE_ONLINE = "/api/fh-admin/skkr/getFuture"; private static String GET_NEW_FUTURE_ONLINE = "/api/fh-admin/skkr/getSkFutureSw"; private static String GET_FUTURE = "/api/fh-admin/skkr/getFuture"; //新水位接口 private static String GET_NEW_FUTURE = "/api/fh-admin/skkr/getSkFutureSw"; //获取是否有预测模型接口 private static String GET_SK_GXZT = "/api/fh-admin/skkr/getSkgxzt"; //同步模型接口 private static String UPDATE_CJYSJ_BY_TM = "/api/fh-admin/skkr/updateCjysjByTm"; // private static String GET_FUTURE_ONLINE = "/api/fh-admin/skkr/getFuture"; // private static String GET_NEW_FUTURE_ONLINE = "/api/fh-admin/skkr/getSkFutureSw"; private static String AUTHORIZATION = "Bearer eyJhbGciOiJSUzI1NiIsInR5cCI6IkpXVCJ9.eyJ1c2VyX25hbWUiOiJBRE1JTiIsIlVzZXJJZCI6IjEiLCJzY29wZSI6WyJhbGwiXSwiVXNlclJlYWxOYW1lIjoi6LaF57qn566h55CG5ZGYIiwiVXNlclh6cWhkbSI6IjQyMDUiLCJleHAiOjI1MTk3NzIwMTAsImp0aSI6ImExNTcxYzk1LTBkNjMtNDkyYi1iOWEyLTE2ZTIzNTQ5ZTY1ZiIsImNsaWVudF9pZCI6InVzZXItc2VydmljZSJ9.ebNVZrw9LbhKaj2w6RR8b2wccQiDkhvBeq79SxxCK-yWiOlIFqBkotTN4TNJg8umcpyYvLILwvqXWRJhffEtgi25sX2y6MqLIWM4kMZ9d8ptdnSmTpBPhltSiQOM0KFa1kl5nSDCBwYOLn-pESJglam76cjpgZNoC88x3iNHacdiXDItY0rtY85HrQ26uyJu9UovKtmYmZRHsIbGMpDta5Q1p4vfaCIr-YUayDrCweZJiQDEEcOSpWJ7O7RMk3pRkX_4UmPHFzrOI2lMp1jQIhxnSTE7EVAz_4z8h8r46muSkpF54Ic4XSawHKqdSLHx8T05LB0MpvOzWMPS6c1uHA"; @@ -37,7 +53,7 @@ HsybGetFuturePo hsybGetFuturePo = apiRequest(GET_FUTURE, resCd); if (hsybGetFuturePo != null && hsybGetFuturePo.getRespCode().equals("200")) { List<List<String>> data = hsybGetFuturePo.getData(); List<List<String>> data = (List<List<String>>)hsybGetFuturePo.getData(); List<List<String>> collect = filterPredictData(data); return collect; @@ -48,29 +64,68 @@ @Override public List<NewSwFuturePo> getNewFuture(String resId) { Map<String, Object> params = new HashMap<>(); params.put("resId", resId); HsybGetFuturePo get = apiRequest(GET_NEW_FUTURE, "GET", params); //获取预测水位之前需要确认是否有预测水位模型 // Boolean isUpdate = isResHasModel(resId); //程序转为hashmap,手动转换类型 List<NewSwFuturePo> data =JSON.parseArray(JSON.toJSONString(get.getData()), NewSwFuturePo.class); List<NewSwFuturePo> filterData = filterPredict(data); // if (isUpdate){ //数据模型已经更新过,就可以去拿预测水位 //设置预测73小时未来水位,因为最后一条是当前时间的整点数据,获取不到三天,多加一小时就可以获取到下一个整点的未来数据 String urlParams = StringUtil.format("?resId={}&yjq=73",resId); log.info(StringUtil.format("过滤后的预测水位数据:{}", JSON.toJSONString(filterData))); return filterData; HsybGetFuturePo hsybGetFuturePo = sendRequestToHsyb(GET_NEW_FUTURE, urlParams); //程序转为hashmap,手动转换类型 List<NewSwFuturePo> data = JSON.parseArray(JSON.toJSONString(hsybGetFuturePo.getData()), NewSwFuturePo.class); List<NewSwFuturePo> filterData = filterPredict(data); log.info(StringUtil.format("过滤后的预测水位数据:{}", JSON.toJSONString(filterData))); return filterData; // }else{ // //拿不到预测水位 // return null; // } } public Boolean isResHasModel(String resId){ String urlParams = StringUtil.format("?skbm={}",resId); HsybGetFuturePo hsybGetFuturePo = sendRequestToHsyb(GET_SK_GXZT, urlParams); String message = hsybGetFuturePo.getData().toString(); if (message.equals("请更新模型数据!")){ return false; }else{ return true; } } @Override public String updateCjysjByTm(List<String> resIds) { // 获取当前时间 LocalDate currentDate = LocalDate.now(); // 获取前30天的时间 LocalDate thirtyDaysBefore = currentDate.minusDays(30); // 输出结果 DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd 00:00:00"); String startTime = thirtyDaysBefore.format(formatter); String endTime = currentDate.format(formatter); String urlParams = StringUtil.format("?startTime={}&endTime={}", startTime, endTime); String skbmListStr = String.join("&skbmList=", resIds); HsybGetFuturePo res = sendRequestToHsyb(UPDATE_CJYSJ_BY_TM, urlParams + "&skbmList=" + skbmListStr); return res.getData().toString(); } private HsybGetFuturePo apiRequest(String url, String type, Map<String, Object> params) { // 获取环境 String activeProfile = SpringContextUtil.getActiveProfile(); if (activeProfile.equals("dev")) { url = LOCAL + url; url = LOCAL + PREFIX + url; } if (activeProfile.equals("prod")) { url = ONLINE + GET_NEW_FUTURE_ONLINE; url = ONLINE + url; } if (activeProfile.equals("test")) { url = ONLINE + GET_NEW_FUTURE_ONLINE; url = ONLINE + url; } HttpMethod method; @@ -99,13 +154,13 @@ // 获取环境 String activeProfile = SpringContextUtil.getActiveProfile(); if (activeProfile.equals("dev")) { url = LOCAL + url; url = LOCAL + PREFIX + url; } if (activeProfile.equals("prod")) { url = ONLINE + GET_FUTURE_ONLINE; url = ONLINE + GET_FUTURE; } if (activeProfile.equals("test")) { url = ONLINE + GET_FUTURE_ONLINE; url = ONLINE + GET_FUTURE; } @@ -131,6 +186,35 @@ System.out.println(e); } return null; } public HsybGetFuturePo sendRequestToHsyb(String url, String urlParam) { // 获取环境 String activeProfile = SpringContextUtil.getActiveProfile(); if (activeProfile.equals("dev")) { url = LOCAL + PREFIX + url; } if (activeProfile.equals("prod")) { url = ONLINE + url; } if (activeProfile.equals("test")) { url = ONLINE + url; } url = url + urlParam; log.info("洪水预报请求地址:{}", url); //设置请求头 HttpHeaders headers = new HttpHeaders(); headers.add("Authorization", AUTHORIZATION); //封装请求头 HttpEntity<MultiValueMap<String, Object>> formEntity = new HttpEntity<MultiValueMap<String, Object>>(headers); try { ResponseEntity<HsybGetFuturePo> exchange = restTemplate.exchange(url, HttpMethod.GET, formEntity, HsybGetFuturePo.class); return exchange.getBody(); } catch (Exception e) { e.printStackTrace(); } return null; } @@ -215,7 +299,7 @@ // return distinctList; } public List<NewSwFuturePo> filterPredict(List<NewSwFuturePo> list){ public List<NewSwFuturePo> filterPredict(List<NewSwFuturePo> list) { /** * { @@ -229,7 +313,7 @@ * "resId": "42092250024" * }, */ if (list.size() == 0){ if (list.size() == 0) { return null; } @@ -240,7 +324,7 @@ List<NewSwFuturePo> filterList = list.stream().filter(item -> item.getTime().indexOf(newSwFuturePo.getTime().split(" ")[1]) > -1).filter(item->!item.getSw().equals("0.0")).collect(Collectors.toList()); //给时间重新赋值,因为南科院hst预测只需要日期 filterList.forEach(e->{ filterList.forEach(e -> { e.setTime(e.getTime().split(" ")[0]); }); skjcmanager/skjcmanager-service/skjcmanager-sm/src/main/java/cn/gistack/sm/sjztmd/feign/AttResBaseClientImpl.java
New file @@ -0,0 +1,30 @@ package cn.gistack.sm.sjztmd.feign; import cn.gistack.sm.sjztmd.entity.AttResBase; import cn.gistack.sm.sjztmd.service.IAttResBaseService; import lombok.AllArgsConstructor; import org.springblade.core.tenant.annotation.NonDS; import org.springblade.core.tool.api.R; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; import springfox.documentation.annotations.ApiIgnore; import java.util.List; import java.util.stream.Collectors; @NonDS @ApiIgnore @RestController @AllArgsConstructor public class AttResBaseClientImpl implements IAttResBaseClient{ private IAttResBaseService attResBaseService; @Override @GetMapping(GET_ALL_RES) public List<String> getAllRes() { List<AttResBase> list = attResBaseService.list(); List<String> collect = list.stream().map(e -> e.getGuid()).collect(Collectors.toList()); return collect; } } skjcmanager/skjcmanager-service/skjcmanager-xxljob/src/main/java/cn/gistack/job/executor/jobhandler/HsybXxlJob.java
New file @@ -0,0 +1,85 @@ package cn.gistack.job.executor.jobhandler; import cn.gistack.nky.fegin.IHsybClient; import cn.gistack.sm.patrol.feign.PatrolTaskClient; import cn.gistack.sm.sjztmd.feign.IAttResBaseClient; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONObject; import com.xxl.job.core.biz.model.ReturnT; import com.xxl.job.core.handler.annotation.XxlJob; import com.xxl.job.core.log.XxlJobLogger; import org.springblade.core.tool.utils.DateUtil; import org.springblade.core.tool.utils.StringUtil; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import java.text.SimpleDateFormat; import java.util.ArrayList; import java.util.List; /** * 洪水预报定时任务执行器 * @author zhongrj * @date 2023-06-03 */ @Component public class HsybXxlJob { @Autowired private IHsybClient hsybClient; @Autowired private IAttResBaseClient attResBaseClient; /** * 洪水预报同步降雨数据 * @param param * @return * @throws Exception */ @XxlJob("updateCjysjByTmHandler") public ReturnT<String> updateCjysjByTmHandler(String param){ XxlJobLogger.log("定时器执行时间:"+ new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(DateUtil.now())); List<String> list = attResBaseClient.getAllRes(); List<List<String>> lists = groupListByQuantity(list, 20); lists.forEach(e->{ XxlJobLogger.log("请求时间:"+ new SimpleDateFormat("HH:mm:ss").format(DateUtil.now())); XxlJobLogger.log("本组水库id为"+e.toString()); XxlJobLogger.log(hsybClient.updateCjysjByTm(e)); XxlJobLogger.log("请求结束时间:"+ new SimpleDateFormat("HH:mm:ss").format(DateUtil.now())); }); XxlJobLogger.log(StringUtil.format("本次共计更新水库数量为{}",list.size())); XxlJobLogger.log("定时器执行结束时间:"+ new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(DateUtil.now())); return ReturnT.SUCCESS; } /** * 将集合按指定数量分组 * * @param list 数据集合 * @param quantity 分组数量 * @return 分组结果 */ public static List<List<String>> groupListByQuantity(List<String> list, int quantity) { if (list == null || list.size() == 0) { return null; } if (quantity <= 0) { throw new IllegalArgumentException("Wrong quantity."); } List<List<String>> wrapList = new ArrayList<List<String>>(); int count = 0; while (count < list.size()) { wrapList.add(new ArrayList<String>(list.subList(count, Math.min((count + quantity), list.size())))); count += quantity; } return wrapList; } }