skjcmanager/skjcmanager-service-api/skjcmanager-alerts-api/src/main/java/cn/gistack/alerts/notice/feign/NoticeClient.java
@@ -18,7 +18,6 @@ /** * 发送告警通知 * @param options 单独超时时间设置 * @param name * @param content * @param type skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/entity/MessageRecord.java
@@ -37,7 +37,7 @@ private String content; /** * 状态 已读,未读 * 状态 0未读 1已读 */ private Integer status; skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/feign/IMessageClient.java
New file @@ -0,0 +1,43 @@ package cn.gistack.sm.message.feign; import cn.gistack.sm.message.entity.MessageRecord; 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; /** * Feign失败配置 * @author zhongrj * @date 2023-07-07 */ @FeignClient( value = "blade-sm", fallback = IMessageClientFallback.class ) public interface IMessageClient { String API_PREFIX = "/client"; String SAVE_MESSAGE_RECORD = API_PREFIX + "/saveMessageRecord"; String GET_NOT_READ_NUMBER = API_PREFIX + "/getNotReadNumber"; /** * 保存站内信记录信息 * @param messageRecord * @return */ @PostMapping(SAVE_MESSAGE_RECORD) void saveMessageRecord(@RequestBody List<MessageRecord> messageRecord); /** * 查询用户未读消息数量 * @param phone * @return */ @GetMapping(GET_NOT_READ_NUMBER) Integer getNotReadNumber(@RequestParam("phone") String phone); } skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/feign/IMessageClientFallback.java
New file @@ -0,0 +1,25 @@ package cn.gistack.sm.message.feign; import cn.gistack.sm.message.entity.MessageRecord; import org.springframework.stereotype.Component; import java.util.List; /** * Feign失败配置 * @author zhongrj * @date 2023-07-07 */ @Component public class IMessageClientFallback implements IMessageClient { @Override public void saveMessageRecord(List<MessageRecord> messageRecord) { } @Override public Integer getNotReadNumber(String phone) { return null; } } skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/vo/MessageRecordVO.java
New file @@ -0,0 +1,13 @@ package cn.gistack.sm.message.vo; import cn.gistack.sm.message.entity.MessageRecord; import lombok.Data; @Data public class MessageRecordVO extends MessageRecord { /** * 是否分页(导出时不分页) */ private Integer isPage; } skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/AsyncNoticeHandle.java
@@ -27,6 +27,38 @@ /** * 水位接近汛限告警 * @param templateId * @param uuid * @param alarmRule * @param adCode */ @Async public void waterComeOverHandle(String templateId, String uuid, AlarmRule alarmRule, String adCode) { // 调用中台服务接口查询数据 JSONArray ztData = noticeHandle.getZtData("", ZtApiUrlConstant.res_diff_dead_z_api); // 只发送站内信 SmsRequestTemplate smsRequestTemplate = noticeHandle.getSendSmsTemplate(ztData, templateId,alarmRule, ZtApiDataColumnConstant.waterLessDeadList,true,null); // 拼接内容后面的联系方式,当前模板县水利部人员姓名电话 // 遍历,拼接县水利部人员电话 for (List<String> list : smsRequestTemplate.getTemplateContent()) { // 后面依次拼接 巡查责任人姓名,技术责任人姓名,市县水利部门人员姓名电话(暂时不要) noticeHandle.addPatrolTechnologyInfo(list); // 市县水利部门人员姓名电话(暂时不要) // list.add("测试:13112341234"); // 保存告警记录信息 // AlarmRecord alarmRecord = noticeHandle.saveAlarmRecord(list,alarmRule,uuid); // if (null != alarmRecord) { // list.add(alarmRecord.getId().toString()); list.add("999"); // } } // 保存站内信记录信息 noticeHandle.saveInStationInfo(smsRequestTemplate,alarmRule); } /** * 水位初次超汛限告警 处理 * @param templateId * @param uuid skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/NoticeHandle.java
@@ -1,6 +1,8 @@ package cn.gistack.alerts.notice.service.impl; import cn.com.flaginfo.sdk.cmc.api.sms.send.SMSSendRequest; import cn.gistack.alerts.alarmRule.entity.AlarmRecord; import cn.gistack.alerts.alarmRule.entity.AlarmRecordDetail; import cn.gistack.alerts.alarmRule.entity.AlarmRule; import cn.gistack.alerts.alarmRule.service.AlarmRecordService; import cn.gistack.alerts.notice.constant.ZtApiUrlConstant; @@ -10,14 +12,19 @@ import cn.gistack.alerts.sms.entity.SmsTemplate; import cn.gistack.alerts.sms.service.ISmsTemplateService; import cn.gistack.alerts.sms.vo.SmsRequestTemplate; import cn.gistack.alerts.sse.server.SSEServer; import cn.gistack.common.utils.SpringContextUtil; import cn.gistack.sm.message.entity.MessageRecord; import cn.gistack.sm.message.feign.IMessageClient; import cn.gistack.sm.sjztmd.feign.IAttResManagePersonClient; import cn.gistack.sm.sjztmd.vo.PersonVO; import cn.gistack.system.feign.ISysClient; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONArray; import com.alibaba.fastjson.JSONObject; import com.alibaba.nacos.shaded.com.google.gson.JsonObject; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import org.springblade.core.secure.utils.AuthUtil; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.HttpEntity; import org.springframework.http.HttpHeaders; @@ -47,6 +54,9 @@ @Autowired private RestTemplate restTemplate; @Autowired private IMessageClient messageClient; @@ -420,4 +430,87 @@ } return null; } /** * 保存站内信记录信息 * @param smsRequestTemplate * @param alarmRule * @return */ public Map<String, Object> saveInStationInfo(SmsRequestTemplate smsRequestTemplate, AlarmRule alarmRule) { // 取出数据 SMSSendRequest request = new SMSSendRequest(); SmsTemplate smsTemplate = smsRequestTemplate.getSmsTemplate(); request.setTemplateId(smsTemplate.getTemplateId()); List<List<String>> templateContent = smsRequestTemplate.getTemplateContent(); // 先保存结果记录 List<AlarmRecordDetail> alarmRecordDetailList = new ArrayList<>(); List<MessageRecord> messageRecordList = new ArrayList<>(); Set<String> phoneSet = new HashSet<>(); // 遍历处理 for (List<String> list : templateContent) { // 设置站内信记录信息 MessageRecord messageRecord = new MessageRecord(); messageRecord.setTheme(alarmRule.getRuleName()); messageRecord.setSender(AuthUtil.getUserId().toString()); messageRecord.setSource("智能告警"); messageRecord.setStatus(2); // 设置告警记录信息 AlarmRecordDetail alarmRecordDetail = new AlarmRecordDetail(); alarmRecordDetail.setReservoirNumber(list.get(0)); alarmRecordDetail.setAlarmRecordId(Long.parseLong(list.get(list.size() - 1))); alarmRecordDetail.setAlarmMode("站内信"); // 未读 alarmRecordDetail.setStatus(1); // 删除 最后一个 告警记录id list.remove(list.size() - 1); // 删除 0 水库编码 list.remove(0); alarmRecordDetail.setPhone(list.get(0)); // 设置手机号 request.setUserNumber(list.get(0)); messageRecord.setRecipient(list.get(0)); phoneSet.add(list.get(0)); // 内容处理 String s = smsTemplate.getContent().replaceAll("\\{.+?\\}", "%s"); // 需要删除最前面的手机号 list.remove(0); // 格式转换 String format = String.format(s, list.toArray()); // 设置内容 request.setMessageContent(format); messageRecord.setContent(format); // 保存记录 alarmRecordDetail.setAlarmContent(format); alarmRecordDetail.setCreateTime(new Date()); // 放入集合 alarmRecordDetailList.add(alarmRecordDetail); messageRecordList.add(messageRecord); } // 保存告警记录详情信息 saveMessageRecord(messageRecordList); Map<String, Object> map = new HashMap<>(2); // 返回 return map; } /** * 保存站内信记录信息 * @param messageRecordList */ public void saveMessageRecord(List<MessageRecord> messageRecordList) { // 保存记录 messageClient.saveMessageRecord(messageRecordList); // 发起通知 // 遍历 for (MessageRecord messageRecord : messageRecordList) { Map<String, Object> map = new HashMap<>(2); // 查询当前用户未读消息个数 map.put("count",messageClient.getNotReadNumber(messageRecord.getRecipient())); map.put("record",messageRecord); SSEServer.sendMessage("web:"+messageRecord.getRecipient(),new JSONObject(map).toJSONString()); } } } skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/NoticeStrategyImpl.java
@@ -196,31 +196,13 @@ * @return */ public String waterComeOverHandle(String arg1, String arg2) { String adCode = null; String templateId = "znxsw10000001"; String uuid = templateId + UUID.randomUUID().toString(); // 查询当前策略对应的告警规则信息 AlarmRule alarmRule = alarmRuleService.getOne(new QueryWrapper<AlarmRule>().eq("rule_name", arg1)); if (null!=alarmRule) { // 调用中台服务接口查询数据 JSONArray ztData = noticeHandle.getZtData("", ZtApiUrlConstant.res_diff_dead_z_api); // 只发送站内信 SmsRequestTemplate smsRequestTemplate = noticeHandle.getSendSmsTemplate(ztData, templateId,alarmRule, ZtApiDataColumnConstant.waterLessDeadList,true,null); // 拼接内容后面的联系方式,当前模板县水利部人员姓名电话 // 遍历,拼接县水利部人员电话 for (List<String> list : smsRequestTemplate.getTemplateContent()) { // 后面依次拼接 巡查责任人姓名,技术责任人姓名,市县水利部门人员姓名电话(暂时不要) noticeHandle.addPatrolTechnologyInfo(list); // 市县水利部门人员姓名电话(暂时不要) // list.add("测试:13112341234"); // 保存告警记录信息 AlarmRecord alarmRecord = noticeHandle.saveAlarmRecord(list,alarmRule,uuid); if (null != alarmRecord) { list.add(alarmRecord.getId().toString()); // list.add("999"); } } // 保存站内信记录信息 // smsService.saveInStationInfo(smsRequestTemplate); asyncNoticeHandle.waterComeOverHandle(templateId,uuid,alarmRule, null); } return "水位接近汛限告警"; } skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/controller/SSEController.java
@@ -1,16 +1,25 @@ package cn.gistack.alerts.sse.controller; import cn.gistack.alerts.notice.service.impl.NoticeHandle; import cn.gistack.alerts.sse.server.SSEServer; import cn.gistack.alerts.sse.vo.SseVO; import cn.gistack.sm.message.entity.MessageRecord; import lombok.AllArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.web.bind.annotation.*; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.util.ArrayList; import java.util.List; @Slf4j @RestController @CrossOrigin @RequestMapping("/sse/sse") @AllArgsConstructor public class SSEController { private final NoticeHandle noticeHandle; /** * 建立连接 @@ -19,7 +28,19 @@ */ @GetMapping("/connect") public SseEmitter connect(SseVO sse){ return SSEServer.connect(sse.getType() + ":" + sse.getUserId()); String userId = sse.getType() + ":" + sse.getUserId(); return SSEServer.connect(userId); } /** * 断开连接 * @param sse * @return */ @GetMapping("/disconnect") public void disconnect(SseVO sse){ String userId = sse.getType() + ":" + sse.getUserId(); SSEServer.removeUser(userId); } /** @@ -27,15 +48,13 @@ * @throws InterruptedException */ @GetMapping("/process") public void sendMessage() throws InterruptedException { SSEServer.sendMessage("web:123","hello!"); // for(int i=0; i<=100; i++){ // if(i>50&&i<70){ // Thread.sleep(500L); // }else{ // Thread.sleep(100L); // } // SSEServer.batchSendMessage(String.valueOf(i)); // } public void sendMessage(){ MessageRecord messageRecord = new MessageRecord(); messageRecord.setRecipient("1642816241683865602"); messageRecord.setTheme("测试123"); messageRecord.setStatus(0); List<MessageRecord> list = new ArrayList<>(); list.add(messageRecord); noticeHandle.saveMessageRecord(list); } } skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/server/SSEServer.java
@@ -34,6 +34,7 @@ //数量+1 count.getAndIncrement(); log.info("create new sse connect ,current user:{}",userId); log.info("count",count.getAndIncrement()); return sseEmitter; } /**