Browse Source

Merge branch 'develop' of com_yunfei_saas/agmp_iots into master

yf_zn 1 year ago
parent
commit
2749f04ea1
35 changed files with 2219 additions and 290 deletions
  1. 15 0
      src/main/java/com/yunfeiyun/agmp/iots/core/mqtt/DeviceTopicService.java
  2. 62 1
      src/main/java/com/yunfeiyun/agmp/iots/device/controller/TestController.java
  3. 6 0
      src/main/java/com/yunfeiyun/agmp/iots/device/service/IRunHaoSfDevice.java
  4. 589 0
      src/main/java/com/yunfeiyun/agmp/iots/device/serviceImp/RunHaoSfDeviceImpl.java
  5. 58 0
      src/main/java/com/yunfeiyun/agmp/iots/domain/IotSfElementfactorAlreadyListResVo.java
  6. 57 0
      src/main/java/com/yunfeiyun/agmp/iots/domain/IotSfElementfactorListReqVo.java
  7. 35 0
      src/main/java/com/yunfeiyun/agmp/iots/domain/IotSfIrrigationRecordListReqVo.java
  8. 41 0
      src/main/java/com/yunfeiyun/agmp/iots/mapper/IotSfElementfactorMapper.java
  9. 27 0
      src/main/java/com/yunfeiyun/agmp/iots/mapper/IotSfIrrigationRecordMapper.java
  10. 32 0
      src/main/java/com/yunfeiyun/agmp/iots/mq/provider/AgmpIotMqProviderService.java
  11. 31 0
      src/main/java/com/yunfeiyun/agmp/iots/mq/provider/AgmpMqProviderService.java
  12. 15 0
      src/main/java/com/yunfeiyun/agmp/iots/service/IIotRunHaoSfdataService.java
  13. 111 0
      src/main/java/com/yunfeiyun/agmp/iots/service/IIotSfElementfactorService.java
  14. 28 0
      src/main/java/com/yunfeiyun/agmp/iots/service/IIotSfIrrigationRecordService.java
  15. 47 0
      src/main/java/com/yunfeiyun/agmp/iots/service/impl/IotRunHaoSfdataServiceImpl.java
  16. 213 0
      src/main/java/com/yunfeiyun/agmp/iots/service/impl/IotSfElementfactorServiceImpl.java
  17. 52 0
      src/main/java/com/yunfeiyun/agmp/iots/service/impl/IotSfIrrigationRecordServiceImpl.java
  18. 13 0
      src/main/java/com/yunfeiyun/agmp/iots/startup/MongoStartup.java
  19. 3 2
      src/main/java/com/yunfeiyun/agmp/iots/task/YbqScheduler.java
  20. 12 4
      src/main/java/com/yunfeiyun/agmp/iots/warn/job/WarnJob.java
  21. 14 4
      src/main/java/com/yunfeiyun/agmp/iots/warn/mapper/IotWarnBussinessMapper.java
  22. 1 1
      src/main/java/com/yunfeiyun/agmp/iots/warn/model/IotWarnconfigCbdInfoVo.java
  23. 10 0
      src/main/java/com/yunfeiyun/agmp/iots/warn/model/WarnResult.java
  24. 64 6
      src/main/java/com/yunfeiyun/agmp/iots/warn/service/IotWarnBussinessService.java
  25. 235 0
      src/main/java/com/yunfeiyun/agmp/iots/warn/service/MsgService.java
  26. 12 2
      src/main/java/com/yunfeiyun/agmp/iots/warn/service/ReCountService.java
  27. 165 0
      src/main/java/com/yunfeiyun/agmp/iots/warn/service/WarnDiseaseService.java
  28. 3 3
      src/main/java/com/yunfeiyun/agmp/iots/warn/service/WarnPestService.java
  29. 3 1
      src/main/java/com/yunfeiyun/agmp/iots/warn/service/WarnService.java
  30. 13 2
      src/main/resources/application-dev.yml
  31. 0 264
      src/main/resources/application-prod.yml
  32. 14 0
      src/main/resources/application-test.yml
  33. 55 0
      src/main/resources/mapper/IotSfElementfactorMapper.xml
  34. 145 0
      src/main/resources/mapper/IotSfIrrigationRecordMapper.xml
  35. 38 0
      src/main/resources/mapper/IotWarnBusinessMapper.xml

+ 15 - 0
src/main/java/com/yunfeiyun/agmp/iots/core/mqtt/DeviceTopicService.java

@@ -98,6 +98,8 @@ public class DeviceTopicService {
                 return getYfXycbIIIBatchSubTopic(deviceId);
             case ServiceNameConst.SERVICE_YF_XCT:
                 return getYfXctDeviceBatchSubTopic(deviceId);
+            case ServiceNameConst.SERVICE_RUNHAO_SF:
+                return getRunHaoSfBatchSubTopic(deviceId);
             default: {
                 throw new IotBizException(IotErrorCode.FAILURE.getCode(), serviceName + "不存在对应topic 解析");
             }
@@ -266,6 +268,19 @@ public class DeviceTopicService {
      * @param deviceId
      * @return
      */
+    private String[] getRunHaoSfBatchSubTopic(String[] deviceId) {
+        String[] topicArray = {
+                IotMqttConstant.RunHaoSfTopic.TOPIC_RUNHAO_SF_REPORT_PREFIX
+        };
+        return getTopics(deviceId, topicArray);
+    }
+
+    /**
+     * 海普发智能开关设备 topic
+     *
+     * @param deviceId
+     * @return
+     */
     private String[] getYrSfDeviceBatchSubTopic(String[] deviceId) {
         String[] topicArray = {
                 IotMqttConstant.YrSfTopic.TOPIC_YR_SF_REPORT_PREFIX,

+ 62 - 1
src/main/java/com/yunfeiyun/agmp/iots/device/controller/TestController.java

@@ -12,12 +12,12 @@ import com.yunfeiyun.agmp.common.utils.uuid.IdUtils;
 import com.yunfeiyun.agmp.iot.common.constant.IotErrorCode;
 import com.yunfeiyun.agmp.iot.common.constant.devicetype.IotDeviceDictConst;
 import com.yunfeiyun.agmp.iot.common.domain.*;
+import com.yunfeiyun.agmp.iot.common.enums.ybq.YbqTypeConst;
 import com.yunfeiyun.agmp.iot.common.exception.IotBizException;
 import com.yunfeiyun.agmp.iot.common.model.cmd.CmdGroupModel;
 import com.yunfeiyun.agmp.iot.common.service.MongoService;
 import com.yunfeiyun.agmp.iot.common.util.dev.DevTypeUtil;
 import com.yunfeiyun.agmp.iots.core.cmd.core.CmdDispatcherService;
-import com.yunfeiyun.agmp.iots.core.manager.ConnectionManager;
 import com.yunfeiyun.agmp.iots.device.domain.WarnTestReq;
 import com.yunfeiyun.agmp.iots.device.mapper.IotDeviceMapper;
 import com.yunfeiyun.agmp.iots.device.service.IYfQxzDevice;
@@ -25,6 +25,7 @@ import com.yunfeiyun.agmp.iots.device.serviceImp.YfQxzDeviceImpl;
 import com.yunfeiyun.agmp.iots.service.IIotDeviceService;
 import com.yunfeiyun.agmp.iots.service.IIotDevicelasteddataService;
 import com.yunfeiyun.agmp.iots.service.IIotYfScddataService;
+import com.yunfeiyun.agmp.iots.service.IotYbqPredictDataService;
 import com.yunfeiyun.agmp.iots.service.impl.IotCbdImgService;
 import com.yunfeiyun.agmp.iots.warn.service.WarnPestService;
 import com.yunfeiyun.agmp.iots.warn.service.WarnService;
@@ -86,6 +87,9 @@ public class TestController {
     private static final DateTimeFormatter stampFormatter = DateTimeFormatter.ofPattern("yyyyMMddHHmmss");
 
 
+    @Autowired
+    private IotYbqPredictDataService iotYbqPredictDataService;
+
     /**
      * 该方法模拟接收到mq的消息解析发送
      *
@@ -765,5 +769,62 @@ public class TestController {
         return AjaxResult.success();
     }
 
+    @RequestMapping("/warn/cmd/ybq")
+    public AjaxResult delaYbqPredictedData(@RequestBody WarnTestReq req) {
+        try {
 
+            JSONObject jobj = req.getData();
+            if (StringUtils.isEmpty(req.getDevCode())) {
+                throw new BizException(ErrorCode.FAILURE.getCode(), "设备code为空");
+            }
+            List<IotDevice> oldIotDevice = iIotDeviceService.selectIotDeviceByDevCode(req.getDevCode());
+            for (IotDevice iotDevice : oldIotDevice) {
+                String code = iotDevice.getDevCode();
+                String devBid = iotDevice.getDevBid();
+                // 得到第三方返回的设备的预测数据
+                String computeDate = jobj.getString("computeDate");
+                String deviceId = jobj.getString("deviceId");
+                String value = jobj.getString("value");
+                //保存预测记录、更新预测最新数据
+                // 查查当天的有没有,有的话更新,没有的话添加
+                IotYbqPredictData iotYbqPredictDataToday = iotYbqPredictDataService.getTodayData(devBid);
+                if (iotYbqPredictDataToday == null) {
+                    IotYbqPredictData iotYbqPredictData = new IotYbqPredictData();
+                    iotYbqPredictData.setId(IdUtils.fastUUID());
+                    iotYbqPredictData.setYbqdataBid(iotYbqPredictData.getId());
+                    iotYbqPredictData.setTid(iotDevice.getTid());
+                    iotYbqPredictData.setComputeDate(computeDate);
+                    iotYbqPredictData.setDateDevType(YbqTypeConst.YBQ_XM_CMB.getCode());
+                    iotYbqPredictData.setValue(value);
+                    iotYbqPredictData.setYbqdataContent(jobj);
+                    iotYbqPredictData.setDeviceId(deviceId);
+                    iotYbqPredictData.setDevBid(devBid);
+                    iotYbqPredictData.setDevTypeBid(iotDevice.getDevtypeBid());
+                    iotYbqPredictData.setYbqdataModifiedDate(DateUtils.dateTimeNow());
+                    iotYbqPredictData.setYbqdataCreatedDate(DateUtils.dateTimeNow());
+                    iotYbqPredictDataService.insertData(iotYbqPredictData);
+                } else {
+                    iotYbqPredictDataToday.setComputeDate(computeDate);
+                    iotYbqPredictDataToday.setValue(value);
+                    iotYbqPredictDataToday.setYbqdataModifiedDate(DateUtils.dateTimeNow());
+                    // 只更新特定的字段
+                    iotYbqPredictDataService.updateData(iotYbqPredictDataToday);
+                }
+                // 同步更新设备信息
+                String extInfo = iotDevice.getExtInfo() == null ? "{}" : iotDevice.getExtInfo();
+                JSONObject jsonObject = JSONObject.parseObject(extInfo);
+                //预测时间
+                jsonObject.put("computeDate", DateUtils.dateTimeNow());
+                //预测value
+                jsonObject.put("computeValue", value);
+                iIotDeviceService.updateIotDeviceExtInfo(devBid, jsonObject.toString());
+                log.info("设备ID={} 的预测数据处理成功,computeDate={}, value={}", iotDevice.getDevBid(), computeDate, value);
+                SpringUtils.getBean(WarnService.class).processWarningDiseaseData();
+            }
+            return AjaxResult.success();
+        } catch (BizException e) {
+            e.printStackTrace();
+        }
+        return AjaxResult.success();
+    }
 }

+ 6 - 0
src/main/java/com/yunfeiyun/agmp/iots/device/service/IRunHaoSfDevice.java

@@ -0,0 +1,6 @@
+package com.yunfeiyun.agmp.iots.device.service;
+
+import com.yunfeiyun.agmp.iots.device.common.Device;
+
+/** 润浩水肥机 */
+public interface IRunHaoSfDevice extends Device {}

+ 589 - 0
src/main/java/com/yunfeiyun/agmp/iots/device/serviceImp/RunHaoSfDeviceImpl.java

@@ -0,0 +1,589 @@
+package com.yunfeiyun.agmp.iots.device.serviceImp;
+
+import com.alibaba.fastjson2.JSONArray;
+import com.alibaba.fastjson2.JSONObject;
+import com.yunfeiyun.agmp.common.utils.DateUtils;
+import com.yunfeiyun.agmp.common.utils.JSONUtils;
+import com.yunfeiyun.agmp.common.utils.StringUtils;
+import com.yunfeiyun.agmp.iot.common.constant.cmd.CmdDef;
+import com.yunfeiyun.agmp.iot.common.constant.devicetype.ServiceNameConst;
+import com.yunfeiyun.agmp.iot.common.constant.mqtt.IotMqttConstant;
+import com.yunfeiyun.agmp.iot.common.domain.IotDevice;
+import com.yunfeiyun.agmp.iot.common.domain.IotRunHaoSfdata;
+import com.yunfeiyun.agmp.iot.common.domain.IotSfElementfactor;
+import com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord;
+import com.yunfeiyun.agmp.iot.common.enums.EnumIrrigationRecord;
+import com.yunfeiyun.agmp.iot.common.enums.EnumSfElementType;
+import com.yunfeiyun.agmp.iot.common.model.cmd.CmdModel;
+import com.yunfeiyun.agmp.iot.common.util.dev.RunHaoSfElementUtil;
+import com.yunfeiyun.agmp.iots.core.manager.MqttManager;
+import com.yunfeiyun.agmp.iots.device.common.DeviceAbstractImpl;
+import com.yunfeiyun.agmp.iots.device.service.IRunHaoSfDevice;
+import com.yunfeiyun.agmp.iots.domain.IotSfElementfactorAlreadyListResVo;
+import com.yunfeiyun.agmp.iots.domain.IotSfElementfactorListReqVo;
+import com.yunfeiyun.agmp.iots.domain.IotSfIrrigationRecordListReqVo;
+import com.yunfeiyun.agmp.iots.service.*;
+import com.yunfeiyun.agmp.iots.service.impl.IotDeviceAddressService;
+import lombok.extern.slf4j.Slf4j;
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.springframework.beans.BeanUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+
+import java.util.*;
+
+/** 润浩水肥机 */
+@Component(ServiceNameConst.SERVICE_RUNHAO_SF)
+@Slf4j
+public class RunHaoSfDeviceImpl extends DeviceAbstractImpl implements IRunHaoSfDevice {
+
+    @Autowired
+    private MqttManager mqttManager;
+
+    @Autowired
+    private IIotDeviceService iIotDeviceService;
+
+    @Autowired
+    private IIotDeviceconfigService iIotDeviceconfigService;
+
+    @Autowired
+    private IIotDevicelasteddataService iIotDevicelasteddataService;
+    @Autowired
+    private IIotRunHaoSfdataService iIotRunHaoSfdataService;
+
+    @Autowired
+    private IIotCmdlogService iIotCmdlogService;
+
+    @Autowired
+    private IotDeviceAddressService iotDeviceAddressService;
+
+    @Autowired
+    private IIotSfElementfactorService iIotSfElementfactorService;
+
+    @Autowired
+    private IIotSfIrrigationRecordService iIotSfIrrigationRecordService;
+
+
+    private void publish(IotDevice iotDevice, String mqttMsgContent) {
+        String devCode = iotDevice.getDevCode();
+        String devBid = iotDevice.getDevBid();
+
+        IotRunHaoSfdata iotRunHaoSfdata = iIotRunHaoSfdataService.selectData(devBid);
+        if (iotRunHaoSfdata != null) {
+            JSONObject msgJson = JSONObject.parseObject(mqttMsgContent);
+            JSONObject jsonObject = iotRunHaoSfdata.getSfdataContent();
+            for(String key: msgJson.keySet()){
+                jsonObject.put(key, msgJson.get(key));
+            }
+            String topic = IotMqttConstant.RunHaoSfTopic.TOPIC_RUNHAO_SF_REPORT_PREFIX + devCode;
+            try{
+                // 延迟1秒 模拟设备上报延迟
+                Thread.sleep(1000);
+                String sendMsg = jsonObject.toString();
+                mqttManager.publishMsg(iotDevice.getDevconnBid(), topic, sendMsg);
+                log.info("【润浩水肥机】发送指令完毕!connectionId:{},topic :{} mqttMsgContent: {}",iotDevice.getDevconnBid(),topic, mqttMsgContent);
+            }catch (Exception e){
+                log.error("【润浩水肥机】发送指令失败!connectionId:{},topic :{} mqttMsgContent: {}",iotDevice.getDevconnBid(),topic, mqttMsgContent);
+            }
+        }
+        //设备暂未对接mqtt,先注释掉
+        // 先手动插入数据库新数据
+
+//        List<String> topicList = new ArrayList<>();
+//        topicList.add(IotMqttConstant.YFScdTopic.TOPIC_SCD_CMD_PREFIX + devCode);
+//        topicList.add(IotMqttConstant.YFScdTopic.TOPIC_SCD_2_CMD_PREFIX + devCode);
+//        if(Objects.equals(iotDevice.getDevtypeBid(), IotDeviceDictConst.TYPE_YF_FXSSCD)){
+//            //新款风吸式杀虫灯
+//            topicList.add(IotMqttConstant.YFScdTopic.TOPIC_FXSSCD_CMD_PREFIX + devCode);
+//        }
+//        for(String topic:topicList){
+//            try{
+//                mqttManager.publishMsg(iotDevice.getDevconnBid(), topic, mqttMsgContent);
+//                log.info("【YFSCD】发送指令完毕!connectionId:{},topic :{} mqttMsgContent: {}",iotDevice.getDevconnBid(),topic, mqttMsgContent);
+//            }catch (Exception e){
+//                log.error("【YFSCD】发送指令失败!connectionId:{},topic :{} mqttMsgContent: {}",iotDevice.getDevconnBid(),topic, mqttMsgContent);
+//            }
+//        }
+    }
+
+    /**
+     * 下发指令
+     *
+     * @param cmdModel
+     * @return
+     */
+    @Override
+    public Object sendCmd(CmdModel cmdModel) throws Exception {
+        log.info("【润浩水肥机】收到指令 任务 cmdModel={}", cmdModel);
+
+        // 获取执行的指令
+        CmdModel.Cmd cmdDistribution = cmdModel.getCmdDistribution();
+        // 获取执行的方法 ,方法可以通过反射获取执行,也可以临时case 匹配
+        String methodName = cmdModel.getCmdDistribution().getFunc();
+
+        String mqttMsgContent = "";
+        String clogSendresult = "发送指令成功";
+        JSONObject jobjParam = null;
+
+        switch (methodName) {
+            case CmdDef.RunHaoSfCmdDef.CMD_GROUP_CONFIG:
+                jobjParam = cmdDistribution.getJsons();
+                String groupCode = jobjParam.getString("sfCode");
+                int groupIndex = Integer.parseInt(groupCode.replace("Btn-qx", ""));
+                String groupIndexStr = String.valueOf(groupIndex);
+                Map<String, String> payloadMap = new HashMap<>();
+                JSONArray valveArray = jobjParam.getJSONArray("childrenList");
+                for (int i = 0; i < valveArray.size(); i++) {
+                    JSONObject valveObj = valveArray.getJSONObject(i);
+                    String valveCode = valveObj.getString("sfCode");
+                    int valveIndex = Integer.parseInt(valveCode.replace("Btn-fa", ""));
+                    String key = "Btn-fx" + String.format("%02d", valveIndex);
+                    payloadMap.put(key, groupIndexStr);
+                }
+
+                mqttMsgContent = JSONUtils.toJSONString(payloadMap);
+                log.info("【润浩水肥机】发送指令【" + CmdDef.RunHaoSfCmdDef.CMD_GROUP_CONFIG + "】 mqttMsgContent={}", mqttMsgContent);
+                break;
+
+            case CmdDef.RunHaoSfCmdDef.CMD_CONFIG:
+                jobjParam = cmdDistribution.getJsons();
+                mqttMsgContent = JSONUtils.toJSONString(jobjParam);
+                log.info("【润浩水肥机】发送指令【config】 mqttMsgContent={}", mqttMsgContent);
+                break;
+
+//            case CmdDef.YfScdCmdDef.CMD_REFRESH:{
+//                JSONObject jobjParam = cmdDistribution.getJsons();
+//                mqttMsgContent = JSONUtils.toJSONString(jobjParam);
+//                log.info("【杀虫灯】发送指令【refresh】 mqttMsgContent={}", mqttMsgContent);
+//                break;
+//            }
+//            case CmdDef.YfScdCmdDef.CMD_COMMON:{
+//                JSONObject jobjParam = cmdDistribution.getJsons();
+//                mqttMsgContent = JSONUtils.toJSONString(jobjParam);
+//                log.info("【杀虫灯】发送指令【report】 mqttMsgContent={}", mqttMsgContent);
+//                break;
+//            }
+        }
+
+
+        if(StringUtils.isNotEmpty(mqttMsgContent)){
+            IotDevice iotDevice= iIotDeviceService.selectIotDeviceByDevBid(cmdModel.getIotDevice().getDevBid());
+            publish(iotDevice, mqttMsgContent);
+        }
+
+        cmdModel.setClogSendresult(clogSendresult);
+        cmdModel.setClogDesc(mqttMsgContent);
+
+        iIotCmdlogService.insertSuccessCmdlog(cmdModel);
+        return null;
+    }
+
+    private IotSfElementfactor getValveElementFactor(IotDevice iotDevice, String valveCode, IotSfElementfactor parentFactor) {
+        IotSfElementfactor valveFactor = RunHaoSfElementUtil.getValveElementFactor(valveCode);
+        if(valveFactor == null){
+            return null;
+        }
+        String sfCreatedDate = parentFactor.getSfCreatedDate();
+        String devBid = iotDevice.getDevBid();
+        String tid = iotDevice.getTid();
+
+        String valveSfBid = valveFactor.getUUId();
+        valveFactor.setSfBid(valveSfBid);
+        valveFactor.setTid(tid);
+        valveFactor.setDevBid(devBid);
+        valveFactor.setSfCreatedDate(sfCreatedDate);
+        valveFactor.setSfModifieddate(sfCreatedDate);
+        valveFactor.setSfParentBid(parentFactor.getSfBid());
+        valveFactor.setSfCreator(iotDevice.getDevCreator());
+        valveFactor.setSfModifier(iotDevice.getDevCreator());
+
+        return valveFactor;
+    }
+
+    private List<IotSfElementfactor> getCreateGroupFactorList(IotDevice iotDevice, String groupCode, List<String> valveList) {
+        List<IotSfElementfactor> createFactorList = new ArrayList<>();
+        IotSfElementfactor groupFactor = RunHaoSfElementUtil.getGroupElementFactor(groupCode);
+        if(groupFactor == null){
+            return createFactorList;
+        }
+
+        String sfCreatedDate = DateUtils.dateTimeNow();
+        String devBid = iotDevice.getDevBid();
+
+        String sfBid = groupFactor.getUUId();
+        groupFactor.setSfBid(sfBid);
+        groupFactor.setTid(iotDevice.getTid());
+        groupFactor.setDevBid(devBid);
+        groupFactor.setSfCreatedDate(sfCreatedDate);
+        groupFactor.setSfModifieddate(sfCreatedDate);
+        groupFactor.setSfCreator(iotDevice.getDevCreator());
+        groupFactor.setSfModifier(iotDevice.getDevCreator());
+
+        createFactorList.add(groupFactor);
+
+        for(String valveCode: valveList){
+            IotSfElementfactor valveFactor = getValveElementFactor(iotDevice, valveCode, groupFactor);
+            if(valveFactor == null){
+                continue;
+            }
+            createFactorList.add(valveFactor);
+        }
+        return createFactorList;
+    }
+
+    public void syncGroupValveConfig(IotDevice iotDevice, JSONObject jsonObject) {
+
+        Map<String, List<String>> groupMap = new HashMap<>();
+        for(String key:jsonObject.keySet()){
+            if(key.startsWith("Btn-fx")){
+                try{
+                    int groupIndex = (int)Math.floor(Double.parseDouble(jsonObject.getString(key)));
+                    if(groupIndex == 0){
+                        continue;
+                    }
+                    String groupCode = String.format("Btn-qx%02d", groupIndex);
+                    int valveIndex = Integer.parseInt(key.replace("Btn-fx", ""));
+                    String valveCode = String.format("Btn-fa%d", valveIndex);
+                    if(!groupMap.containsKey(groupCode)){
+                        groupMap.put(groupCode, new ArrayList<>());
+                    }
+                    groupMap.get(groupCode).add(valveCode);
+                }catch (Exception e){
+                    continue;
+                }
+            }
+        }
+
+        String devBid = iotDevice.getDevBid();
+        IotSfElementfactorListReqVo reqVo = new IotSfElementfactorListReqVo();
+        reqVo.setDevBid(devBid);
+        reqVo.setTid(iotDevice.getTid());
+        List<IotSfElementfactorAlreadyListResVo> elementList = iIotSfElementfactorService.getGroupAlreadyElementList(reqVo);
+        Map<String, IotSfElementfactorAlreadyListResVo> eleMap = new HashMap<>();
+        for(IotSfElementfactorAlreadyListResVo ele: elementList){
+            eleMap.put(ele.getSfCode(), ele);
+        }
+
+        List<IotSfElementfactor> createFactorList = new ArrayList<>();
+        String sfCreatedDate = DateUtils.dateTimeNow();
+        List<String> deleteSfBidList = new ArrayList<>();
+        for(Map.Entry<String, List<String>> entry: groupMap.entrySet()){
+            String groupCode = entry.getKey();
+            List<String> valveList = entry.getValue();
+            // 先删除已存在的灌区配置,如果最后灌区仍有剩余,则表示设备已经删除,平台未删除,需要删除平台的 与设备同步
+            IotSfElementfactorAlreadyListResVo groupEle = eleMap.remove(groupCode);
+            // 如果不存在,则创建新的灌区配置
+            if(groupEle == null){
+                List<IotSfElementfactor> createGroupList = getCreateGroupFactorList(iotDevice, groupCode, valveList);
+                createFactorList.addAll(createGroupList);
+            }else{
+                // 如果存在,则比较灌区内的电磁阀配置,如果不同,则删除原有的电磁阀配置,创建新的电磁阀配置
+                List<IotSfElementfactorAlreadyListResVo> valveEleList = groupEle.getChildrenList();
+                if(valveEleList == null){
+                    valveEleList = new ArrayList<>();
+                }
+                Map<String, IotSfElementfactorAlreadyListResVo> valveEleMap = new HashMap<>();
+                for(IotSfElementfactorAlreadyListResVo valveEle: valveEleList){
+                    valveEleMap.put(valveEle.getSfCode(), valveEle);
+                }
+                for(String valveCode: valveList){
+                    IotSfElementfactorAlreadyListResVo valveEle = valveEleMap.remove(valveCode);
+                    // 如果存在,则不处理
+                    // 如果不存在,则创建当前灌区下的新电磁阀配置
+                    if(valveEle == null){
+                        IotSfElementfactor parentFactor = new IotSfElementfactor();
+                        BeanUtils.copyProperties(groupEle, parentFactor);
+                        parentFactor.setSfCreatedDate(sfCreatedDate);
+
+                        IotSfElementfactor valveFactor = getValveElementFactor(iotDevice, valveCode, parentFactor);
+                        if(valveFactor == null){
+                            continue;
+                        }
+                        createFactorList.add(valveFactor);
+                    }
+                }
+                // 剩余的元素为需要删除的元素
+                for(IotSfElementfactorAlreadyListResVo valveEle: valveEleMap.values()){
+                    deleteSfBidList.add(valveEle.getSfBid());
+                }
+            }
+        }
+
+        // 剩余的元素为需要删除的元素
+        for(IotSfElementfactorAlreadyListResVo ele: eleMap.values()){
+            if(ele.getChildrenList() != null && !ele.getChildrenList().isEmpty()){
+                deleteSfBidList.add(ele.getSfBid());
+            }
+        }
+        if(!deleteSfBidList.isEmpty()){
+            iIotSfElementfactorService.batchDeleteIotSfElementfactorBySfBidList(deleteSfBidList);
+        }
+        // 创建新的灌区配置
+        if(!createFactorList.isEmpty()){
+            iIotSfElementfactorService.batchInsertIotSfElementfactor(createFactorList);
+        }
+    }
+
+    public void syncIrrigationRecord(IotDevice iotDevice, JSONObject jsonObject){
+        String runMode = "Btn-zdsd";
+        String runStatus = "Btn-yjqd";
+
+        String zdsd = jsonObject.getString(runMode);
+        String yjqd = jsonObject.getString(runStatus);
+
+        // 如果进行中,则检测是否有未完成的灌溉记录,如果有,则不处理,如果没有,则创建新的灌溉记录
+        if("1".equals(yjqd)){
+            // 暂不处理
+            return;
+        }
+        IotSfElementfactorListReqVo reqVo = new IotSfElementfactorListReqVo();
+        reqVo.setDevBid(iotDevice.getDevBid());
+        reqVo.setTid(iotDevice.getTid());
+        reqVo.setSfType(EnumSfElementType.SUCTION.getCode());
+        List<IotSfElementfactor> elementfactorList = iIotSfElementfactorService.selectIotSfElementfactorList(reqVo);
+        if(elementfactorList.isEmpty()){
+            return;
+        }
+
+        Map<IotSfElementfactor, Double> llMap = new HashMap<>();
+
+        for(IotSfElementfactor elementFactor: elementfactorList) {
+            String sfCode = elementFactor.getSfCode();
+            String key = sfCode.replace("Btn-fs", "Num-ls");
+            String v = jsonObject.getString(key);
+            Double ll = 0.0;
+            try{
+                ll = Double.parseDouble(v);
+            }catch (Exception e){}
+            llMap.put(elementFactor, ll);
+        }
+
+        IotSfIrrigationRecordListReqVo reqVo1 = new IotSfIrrigationRecordListReqVo();
+        reqVo1.setDevBid(iotDevice.getDevBid());
+        reqVo1.setTid(iotDevice.getTid());
+        reqVo1.setRcdMode(zdsd);
+        reqVo1.setRcdStatus(EnumIrrigationRecord.STATUS_RUNNING.getCode());
+
+        List<IotSfIrrigationRecord> recordList = iIotSfIrrigationRecordService.selectIrrigationRecordList(reqVo1);
+        if(recordList.isEmpty()){
+            return;
+        }
+        Set<String> sfdataBidSet = new HashSet<>();
+        for(IotSfIrrigationRecord record: recordList){
+            sfdataBidSet.add(record.getSfdataBid());
+        }
+        List<IotRunHaoSfdata> iotRunHaoSfdataList = iIotRunHaoSfdataService.selectDataList(new ArrayList<>(sfdataBidSet));
+        if(iotRunHaoSfdataList.isEmpty()){
+            return;
+        }
+        Map<String, IotRunHaoSfdata> iotRunHaoSfdataMap = new HashMap<>();
+        for(IotRunHaoSfdata iotRunHaoSfdata: iotRunHaoSfdataList){
+            iotRunHaoSfdataMap.put(iotRunHaoSfdata.getSfdataBid(), iotRunHaoSfdata);
+        }
+        String endDateStr = DateUtils.dateTimeNow();
+        Date endDate = DateUtils.parseDate(endDateStr);
+        for(IotSfIrrigationRecord record: recordList){
+            IotRunHaoSfdata iotRunHaoSfdata = iotRunHaoSfdataMap.get(record.getSfdataBid());
+            JSONObject sfdataContent = iotRunHaoSfdata.getSfdataContent();
+            double totalLL = 0.0;
+            for(Map.Entry<IotSfElementfactor, Double> entry: llMap.entrySet()) {
+                IotSfElementfactor elementFactor = entry.getKey();
+                String sfCode = elementFactor.getSfCode();
+                String key = sfCode.replace("Btn-fs", "Num-ls");
+
+                Double oldLL = entry.getValue();
+                double newLL = 0.0;
+                try{
+                    newLL = Double.parseDouble(sfdataContent.getString(key));
+                }catch (Exception e){}
+
+                double diffLL = newLL - oldLL;
+                if(diffLL < 0){
+                    diffLL = 0.0;
+                }
+                llMap.put(elementFactor, diffLL);
+                totalLL += diffLL;
+            }
+            StringBuilder rcdContent = new StringBuilder(record.getRcdContent());
+            rcdContent = new StringBuilder(rcdContent.toString().replace(EnumIrrigationRecord.STATUS_RUNNING.getName(), " 灌溉完成. "));
+            String rcdFertilizer = "";
+            for(Map.Entry<IotSfElementfactor, Double> entry: llMap.entrySet()) {
+                IotSfElementfactor elementFactor = entry.getKey();
+                String sfDisplayname = elementFactor.getSfDisplayname();
+                Double diffLL = entry.getValue();
+                String msg = sfDisplayname + "用量:" + diffLL + "L  ";
+                rcdFertilizer += msg;
+            }
+
+            Date startDate = DateUtils.parseDate(record.getRcdStartdate());
+            long diff = endDate.getTime() - startDate.getTime();
+            int diffMinutes = (int) (diff / (60 * 1000));
+            if(diffMinutes < 0){
+                diffMinutes = 0;
+            }
+
+            record.setRcdFertilizer(rcdFertilizer);
+            record.setRcdTime(diffMinutes);
+            record.setRcdContent(rcdContent.toString());
+            record.setRcdStatus(EnumIrrigationRecord.STATUS_FINISHED.getCode());
+            record.setRcdFlow(totalLL);
+            record.setRcdEnddate(endDateStr);
+        }
+        iIotSfIrrigationRecordService.batchUpdateIrrigationRecord(recordList);
+    }
+
+
+    public Object cmdData(JSONObject dataJson, String topic, String connectionId, String devUpdateddate) throws Exception {
+        log.info("润浩水肥 数据解析 {},topic:{}", dataJson.toString(),topic);
+
+        IotDevice oldIotDevice = findIotDevice(topic, dataJson, connectionId);
+        if (oldIotDevice == null) {
+            log.error("未取到 iotDevice");
+            return null;
+        }
+
+        IotDevice iotDevice = new IotDevice();
+        iotDevice.setTid(oldIotDevice.getTid());
+        iotDevice.setDevtypeBid(oldIotDevice.getDevtypeBid());
+        iotDevice.setDevBid(oldIotDevice.getDevBid());
+        iotDevice.setDevUpdateddate(devUpdateddate);
+        iotDevice.setDevStatus("1");//设备上线
+
+        String[] keyArrays = {
+                "Btn-dsdl",   // 施肥模式  1 定时模式 0 定量模式
+                "Btn-zfxz",   // 注肥开关  1 开 0 关
+                "Btn-jbms",   // 搅拌模式  1 联动模式 0 搅拌模式
+                "Num-jbsjA",  // 搅拌时间  A搅拌机 单位 分钟
+                "Num-jbsjB",
+                "Num-jbsjC",
+                "Num-jbsjD",
+                "Btn-jbA",    // 搅拌开关  A搅拌机   1 开 0 关
+                "Btn-jbB",
+                "Btn-jbC",
+                "Btn-jbD",
+                "Num-jbsyA",  // 剩余搅拌时间  A搅拌机 单位 分钟
+                "Num-jbsyB",
+                "Num-jbsyC",
+                "Num-jbsyD",
+                "Num-lgcs",   // 轮灌次数  单位 次  0 - 10
+                "Num-lgjg",   // 轮灌间隔  单位 分钟  0 - 1000
+                "Btn-zdsd"    // 运行模式  1 自动 0 手动
+        };
+//
+        JSONObject extConf = new JSONObject();
+        for (String k : keyArrays) {
+            String v = "0";
+            if (dataJson.containsKey(k)) {
+                v = dataJson.getString(k);
+                if("true".equals(v)){
+                    v = "1";
+                }else if("false".equals(v)){
+                    v = "0";
+                }
+                v = String.valueOf((int) Math.floor(Double.parseDouble(v)));
+            }
+            extConf.put(k, v);
+        }
+        String devConfig = JSONUtils.toJSONString(extConf);
+
+        // 更新设备基础信息数据库 mysql
+        iIotDeviceService.updateIotDevice(iotDevice);
+        // 创建或更新设备配置信息
+        if (StringUtils.isNotEmpty(devConfig)) {
+            iIotDeviceconfigService.createOrUpdateDevConfig(oldIotDevice, devConfig, iotDevice.getDevUpdateddate());
+        }
+
+        // 更新设备数据到mongodb
+        iIotRunHaoSfdataService.insertData(iotDevice, dataJson);
+
+        // 保存 设备最新数据 到redis
+        iIotDevicelasteddataService.updateDeviceLastedData(oldIotDevice, String.valueOf(dataJson), devUpdateddate);
+
+        try{
+            // 同步灌区配置
+            syncGroupValveConfig(oldIotDevice, dataJson);
+        }catch (Exception e){
+            log.error("润浩水肥机 同步灌区配置失败", e);
+        }
+        try{
+            // 同步灌溉记录
+            syncIrrigationRecord(oldIotDevice, dataJson);
+        }catch (Exception e){
+            log.error("润浩水肥机 同步灌溉记录失败", e);
+        }
+        return null;
+    }
+
+    public void cmdOffline(JSONObject dataJson, String topic, String connectionId) throws MqttException {
+        log.debug("杀虫灯离线数据 {}", dataJson.toString());
+//        IotDevice oldIotDevice = findIotDevice(topic, dataJson, connectionId);
+//        if (oldIotDevice == null) {
+//            log.error("未取到 iotDevice");
+//            return;
+//        }
+//
+//        IotDevice newIotDevice = new IotDevice();
+//        newIotDevice.setDevBid(oldIotDevice.getDevBid());
+//        newIotDevice.setDevStatus("0");
+//        newIotDevice.setDevOfflinedate(DateUtils.dateTimeNow());
+//        newIotDevice.setTid(oldIotDevice.getTid());
+//        newIotDevice.setDevtypeBid(oldIotDevice.getDevtypeBid());
+//        newIotDevice.setDevCreateddate(oldIotDevice.getDevCreateddate());
+//        newIotDevice.setDevUpdateddate(oldIotDevice.getDevUpdateddate());
+//        newIotDevice.setDevCode(oldIotDevice.getDevCode());
+//        newIotDevice.setDevOriginalStatus(oldIotDevice.getDevOriginalStatus());
+//        iIotDeviceService.updateIotDevice(newIotDevice);
+//
+//
+//        //发送离线预警
+//        SpringUtils.getBean(WarnService.class).processWarningOfflineData(newIotDevice, dataJson);
+//
+//        /**
+//         * 下发刷新指令,检测设备是否真离线
+//         */
+//        JSONObject payload = new JSONObject();
+//        payload.put("cmd", "read");
+//        payload.put("ext", "data");
+//        String mqttMsgContent = JSONUtils.toJSONString(payload);
+//        publish(oldIotDevice, mqttMsgContent);
+//
+//        log.info("[杀虫灯] 下发刷新指令,检测设备是否真离线: " + oldIotDevice.getDevCode());
+    }
+
+    /**
+     * 接收上报数据
+     *
+     * @param topic
+     * @param dataJson
+     * @return
+     */
+    @Override
+    public Object receiveData(String topic, JSONObject dataJson,String connectionId) throws Exception {
+        log.info("润浩水肥机实现类  处理收到的 设备上报数据 " + dataJson.toString());
+        // 接收设备上报数据后的处理逻辑
+        String devUpdateddate = dataJson.getString("devUpdateddate");
+        if(StringUtils.isEmpty(devUpdateddate)){
+            devUpdateddate= DateUtils.dateTimeNow();
+        }
+        this.cmdData(dataJson, topic, connectionId, devUpdateddate);
+        return null;
+    }
+
+    @Override
+    public boolean isDeviceProps(JSONObject jobjMsg) {
+        return "data".equals(jobjMsg.getString("cmd"));
+    }
+
+    /**
+     * 根据topic、设备发来的消息,查询对应设备实体
+     *
+     * @param topic
+     * @param jobjMsg
+     * @return
+     */
+    @Override
+    public IotDevice findIotDevice(String topic, JSONObject jobjMsg,String connectionId) {
+        String devId = mqttManager.getDevIdByTopic(connectionId,topic);
+        return iIotDeviceService.selectIotDeviceByDevBid(devId);
+    }
+}

+ 58 - 0
src/main/java/com/yunfeiyun/agmp/iots/domain/IotSfElementfactorAlreadyListResVo.java

@@ -0,0 +1,58 @@
+package com.yunfeiyun.agmp.iots.domain;
+
+import com.yunfeiyun.agmp.iot.common.domain.IotBaseEntity;
+import lombok.Data;
+
+import java.util.List;
+
+/**
+ * 水肥机要素
+ */
+@Data
+public class IotSfElementfactorAlreadyListResVo extends IotBaseEntity {
+    private static final long serialVersionUID = 1L;
+
+    /** 自增主键 */
+    private Long id;
+
+    /** 业务标识 */
+    private String sfBid;
+
+    /** 设备编号 */
+    private String devBid;
+
+    /** 水肥要素类型 */
+    private String sfType;
+
+    /** 原始名称 */
+    private String sfName;
+
+    /** 显示名称 */
+    private String sfDisplayname;
+
+    /** 要素编码 */
+    private String sfCode;
+
+    /** 父类id */
+    private String sfParentBid;
+
+    /** 排序字段 默认0 */
+    private Integer sfSequence;
+
+    /** 租户ID */
+    private String tid;
+
+    /** 创建时间 */
+    private String sfCreatedDate;
+
+    /** 创建人 */
+    private String sfCreator;
+
+    /** 修改时间 */
+    private String sfModifieddate;
+    /** 修改人 */
+    private String sfModifier;
+
+    private List<IotSfElementfactorAlreadyListResVo> childrenList;
+
+}

+ 57 - 0
src/main/java/com/yunfeiyun/agmp/iots/domain/IotSfElementfactorListReqVo.java

@@ -0,0 +1,57 @@
+package com.yunfeiyun.agmp.iots.domain;
+
+import lombok.Data;
+
+import java.util.List;
+
+/**
+ * 水肥机要素
+ */
+@Data
+public class IotSfElementfactorListReqVo {
+    private static final long serialVersionUID = 1L;
+
+    /** 自增主键 */
+    private Long id;
+
+    /** 业务标识 */
+    private String sfBid;
+
+    /** 设备编号 */
+    private String devBid;
+
+    /** 水肥要素类型 */
+    private String sfType;
+
+    /** 原始名称 */
+    private String sfName;
+
+    /** 显示名称 */
+    private String sfDisplayname;
+
+    /** 要素编码 */
+    private String sfCode;
+
+    /** 父类id */
+    private String sfParentBid;
+
+    /** 排序字段 默认0 */
+    private Integer sfSequence;
+
+    /** 租户ID */
+    private String tid;
+
+    /** 创建时间 */
+    private String sfCreatedDate;
+
+    /** 创建人 */
+    private String sfCreator;
+
+    /** 修改时间 */
+    private String sfModifieddate;
+    /** 修改人 */
+    private String sfModifier;
+
+    private List<String> sfTypeList;
+
+}

+ 35 - 0
src/main/java/com/yunfeiyun/agmp/iots/domain/IotSfIrrigationRecordListReqVo.java

@@ -0,0 +1,35 @@
+package com.yunfeiyun.agmp.iots.domain;
+
+import com.yunfeiyun.agmp.iot.common.domain.IotBaseEntity;
+import lombok.Data;
+
+@Data
+public class IotSfIrrigationRecordListReqVo extends IotBaseEntity {
+
+    // 设备标识
+    private String devBid;
+
+    // 灌区标识
+    private String rcdGroupbid;
+
+    // 灌区名称
+    private String rcdGroupName;
+
+    // 灌溉状态 0 进行中 1 已完成
+    private String rcdStatus;
+
+    // 灌溉模式 0 手动 1 自动
+    private String rcdMode;
+
+    // 灌溉数据标识
+    private String sfdataBid;
+
+    // 开始时间
+    private String startTime;
+
+    // 结束时间
+    private String endTime;
+
+    private String tid;
+
+}

+ 41 - 0
src/main/java/com/yunfeiyun/agmp/iots/mapper/IotSfElementfactorMapper.java

@@ -0,0 +1,41 @@
+package com.yunfeiyun.agmp.iots.mapper;
+
+import com.yunfeiyun.agmp.iot.common.domain.IotSfElementfactor;
+import com.yunfeiyun.agmp.iots.domain.IotSfElementfactorListReqVo;
+
+import java.util.List;
+
+/**
+ * Mapper接口
+ *
+ */
+public interface IotSfElementfactorMapper {
+
+    /**
+     * 查询水肥机要素列表
+     *
+     * @param iotSfElementfactor 水肥机要素
+     * @return 水肥机要素集合
+     */
+    public List<IotSfElementfactor> selectIotSfElementfactorList(IotSfElementfactorListReqVo reqVo);
+
+    /**
+     * 批量插入水肥机要素
+     *
+     * @param factorList 水肥机要素列表
+     * @return 结果
+     */
+    public int batchInsertIotSfElementfactor(List<IotSfElementfactor> factorList);
+
+
+    /**
+     * 删除水肥机要素
+     *
+     * @param sfBid 水肥机要素ID
+     * @return 结果
+     */
+    public int deleteIotSfElementfactorBySfBid(String sfBid);
+
+    public int batchDeleteIotSfElementfactorBySfBidList(List<String> sfBidList);
+
+}

+ 27 - 0
src/main/java/com/yunfeiyun/agmp/iots/mapper/IotSfIrrigationRecordMapper.java

@@ -0,0 +1,27 @@
+package com.yunfeiyun.agmp.iots.mapper;
+
+import com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord;
+import com.yunfeiyun.agmp.iots.domain.IotSfIrrigationRecordListReqVo;
+
+import java.util.List;
+
+public interface IotSfIrrigationRecordMapper {
+    // 添加灌溉记录
+    public int insertIrrigationRecord(IotSfIrrigationRecord record);
+
+    public int batchInsertIotSfIrrigationRecord(List<IotSfIrrigationRecord> recordList);
+
+    // 根据唯一标识查询灌溉记录
+    public IotSfIrrigationRecord selectIrrigationRecordByBid(String rcdBid);
+
+    // 更新灌溉记录
+    public int updateIrrigationRecord(IotSfIrrigationRecord record);
+
+//    // 删除灌溉记录
+//    public void deleteIrrigationRecord(String rcdBid);
+
+    // 查询所有灌溉记录
+    public List<IotSfIrrigationRecord> selectIrrigationRecordList(IotSfIrrigationRecordListReqVo record);
+
+    public int batchUpdateIrrigationRecord(List<IotSfIrrigationRecord> recordList);
+}

+ 32 - 0
src/main/java/com/yunfeiyun/agmp/iots/mq/provider/AgmpIotMqProviderService.java

@@ -0,0 +1,32 @@
+package com.yunfeiyun.agmp.iots.mq.provider;
+
+import com.yunfeiyun.agmp.common.framework.mq.rabbitmq.consts.MqAgmpConsts;
+import com.yunfeiyun.agmp.common.framework.mq.rabbitmq.consts.MqAgmpIotConsts;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
+import org.springframework.stereotype.Service;
+
+@Service
+@Slf4j
+@ConditionalOnBean(name = "agmpIotMqConfig")
+public class AgmpIotMqProviderService {
+    @Autowired
+    @Qualifier("agmpIotRabbitTemplate")
+    private RabbitTemplate agmpIotRabbitTemplate;
+
+
+    /**
+     * 往智慧农业发送
+     *
+     * @param message
+     */
+    public void sendToAgmpIot(String message) {
+        log.info("【消息通知】sendToagmpIot message:{}", message);
+        agmpIotRabbitTemplate.convertAndSend(MqAgmpIotConsts.ExchangeConsts.EXCHANGE_KEY, MqAgmpIotConsts.RoutingConsts.ROUTING_KEY, message);
+    }
+
+
+}

+ 31 - 0
src/main/java/com/yunfeiyun/agmp/iots/mq/provider/AgmpMqProviderService.java

@@ -0,0 +1,31 @@
+package com.yunfeiyun.agmp.iots.mq.provider;
+
+import com.yunfeiyun.agmp.common.framework.mq.rabbitmq.consts.MqAgmpConsts;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
+import org.springframework.stereotype.Service;
+
+@Service
+@Slf4j
+@ConditionalOnBean(name = "agmpMqConfig")
+public class AgmpMqProviderService {
+    @Autowired
+    @Qualifier("agmpRabbitTemplate")
+    private RabbitTemplate agmpRabbitTemplate;
+
+
+    /**
+     * 往智慧农业发送
+     *
+     * @param message
+     */
+    public void sendToAgmp(String message) {
+        log.info("【消息通知】sendToAgmp message:{}", message);
+        agmpRabbitTemplate.convertAndSend(MqAgmpConsts.ExchangeConsts.EXCHANGE_KEY, MqAgmpConsts.RoutingConsts.ROUTING_KEY, message);
+    }
+
+
+}

+ 15 - 0
src/main/java/com/yunfeiyun/agmp/iots/service/IIotRunHaoSfdataService.java

@@ -0,0 +1,15 @@
+package com.yunfeiyun.agmp.iots.service;
+
+import com.alibaba.fastjson2.JSONObject;
+import com.yunfeiyun.agmp.iot.common.domain.IotDevice;
+import com.yunfeiyun.agmp.iot.common.domain.IotRunHaoSfdata;
+
+import java.util.List;
+
+public interface IIotRunHaoSfdataService {
+    public void insertData(IotDevice iotDevice, JSONObject jsonObject);
+
+    public IotRunHaoSfdata selectData(String devBid);
+
+    public List<IotRunHaoSfdata> selectDataList(List<String> sfdataBidList);
+}

+ 111 - 0
src/main/java/com/yunfeiyun/agmp/iots/service/IIotSfElementfactorService.java

@@ -0,0 +1,111 @@
+package com.yunfeiyun.agmp.iots.service;
+
+import com.yunfeiyun.agmp.iot.common.domain.IotSfElementfactor;
+import com.yunfeiyun.agmp.iots.domain.IotSfElementfactorAlreadyListResVo;
+import com.yunfeiyun.agmp.iots.domain.IotSfElementfactorListReqVo;
+
+import java.util.List;
+
+/**
+ * 水肥机要素Service
+ */
+public interface IIotSfElementfactorService {
+    /**
+     * 查询水肥机要素列表
+     *
+     * @param iotSfElementfactor 水肥机要素
+     * @return 水肥机要素集合
+     */
+    public List<IotSfElementfactor> selectIotSfElementfactorList(IotSfElementfactorListReqVo reqVo);
+
+    /**
+     * 查询泵类已配置原始要素列表
+     * @param reqVo
+     * @return
+     */
+
+    public List<IotSfElementfactor> selectIotSfElementfactorListByPump(IotSfElementfactorListReqVo reqVo);
+
+    /**
+     * 查询灌区已配置原始要素列表
+     * @param reqVo
+     * @return
+     */
+
+    public List<IotSfElementfactor> selectIotSfElementfactorListByGroup(IotSfElementfactorListReqVo reqVo);
+
+    /**
+     * 查询电磁阀已配置原始要素列表
+     * @param reqVo
+     * @return
+     */
+
+    public List<IotSfElementfactor> selectIotSfElementfactorListByValve(IotSfElementfactorListReqVo reqVo);
+
+
+
+    /**
+     * 批量插入水肥机要素
+     *
+     * @param factorList 水肥机要素列表
+     * @return 结果
+     */
+    public int batchInsertIotSfElementfactor(List<IotSfElementfactor> factorList);
+
+    /**
+     * 删除水肥机要素
+     *
+     * @param sfBid 水肥机要素ID
+     * @return 结果
+     */
+    public int deleteIotSfElementfactorBySfBid(String sfBid);
+
+    /**
+     * 批量删除水肥机要素
+     *
+     * @param sfBid 水肥机要素ID
+     * @return 结果
+     */
+    public int batchDeleteIotSfElementfactorBySfBidList(List<String> sfBidList);
+
+//    /**
+//     * 查询已配置要素列表
+//     * @param reqVo
+//     * @return
+//     */
+//
+//    public List<IotSfElementfactorAlreadyListResVo> getAlreadyElementList(IotSfElementfactorListReqVo reqVo);
+
+    /**
+     * 查询已配置要素列表
+     * @param reqVo
+     * @return
+     */
+
+    public List<IotSfElementfactorAlreadyListResVo> getAlreadyElementList(List<IotSfElementfactor> factorList);
+
+    /**
+     * 查询泵类已配置要素列表
+     * @param reqVo
+     * @return
+     */
+
+    public List<IotSfElementfactorAlreadyListResVo> getPumpAlreadyElementList(IotSfElementfactorListReqVo reqVo);
+
+    /**
+     * 查询灌区已配置要素列表
+     * @param reqVo
+     * @return
+     */
+
+    public List<IotSfElementfactorAlreadyListResVo> getGroupAlreadyElementList(IotSfElementfactorListReqVo reqVo);
+
+    /**
+     * 查询电磁阀已配置要素列表
+     * @param reqVo
+     * @return
+     */
+
+    public List<IotSfElementfactorAlreadyListResVo> getValveAlreadyElementList(IotSfElementfactorListReqVo reqVo);
+}
+

+ 28 - 0
src/main/java/com/yunfeiyun/agmp/iots/service/IIotSfIrrigationRecordService.java

@@ -0,0 +1,28 @@
+package com.yunfeiyun.agmp.iots.service;
+
+import com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord;
+import com.yunfeiyun.agmp.iots.domain.IotSfIrrigationRecordListReqVo;
+
+import java.util.List;
+
+public interface IIotSfIrrigationRecordService {
+    // 添加灌溉记录
+    public int insertIrrigationRecord(IotSfIrrigationRecord record);
+
+    public int batchInsertIotSfIrrigationRecord(List<IotSfIrrigationRecord> recordList);
+
+    // 根据唯一标识查询灌溉记录
+    public IotSfIrrigationRecord selectIrrigationRecordByBid(String rcdBid);
+
+    // 更新灌溉记录
+    public int updateIrrigationRecord(IotSfIrrigationRecord record);
+
+    // 更新灌溉记录
+    public int batchUpdateIrrigationRecord(List<IotSfIrrigationRecord> recordList);
+
+    // 删除灌溉记录
+    public void deleteIrrigationRecord(String rcdBid);
+
+    // 查询所有灌溉记录
+    public List<IotSfIrrigationRecord> selectIrrigationRecordList(IotSfIrrigationRecordListReqVo record);
+}

+ 47 - 0
src/main/java/com/yunfeiyun/agmp/iots/service/impl/IotRunHaoSfdataServiceImpl.java

@@ -0,0 +1,47 @@
+package com.yunfeiyun.agmp.iots.service.impl;
+
+import com.alibaba.fastjson2.JSONObject;
+import com.yunfeiyun.agmp.iot.common.domain.IotDevice;
+import com.yunfeiyun.agmp.iot.common.domain.IotRunHaoSfdata;
+import com.yunfeiyun.agmp.iot.common.service.MongoService;
+import com.yunfeiyun.agmp.iots.service.IIotRunHaoSfdataService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+@Service
+@Slf4j
+public class IotRunHaoSfdataServiceImpl implements IIotRunHaoSfdataService {
+    @Autowired
+    private MongoService mongoService;
+
+    @Override
+    public void insertData(IotDevice iotDevice, JSONObject jsonObject) {
+        IotRunHaoSfdata iotRunHaoSfdata = new IotRunHaoSfdata();
+        iotRunHaoSfdata.setCId(iotDevice.getTid());
+        iotRunHaoSfdata.setSfdataBid(iotRunHaoSfdata.getUUId());
+        iotRunHaoSfdata.setDevBid(iotDevice.getDevBid());
+        iotRunHaoSfdata.setSfdataCreatedDate(iotDevice.getDevUpdateddate());
+        iotRunHaoSfdata.setSfdataContent(jsonObject);
+
+        mongoService.saveOne(iotRunHaoSfdata);
+    }
+
+    @Override
+    public IotRunHaoSfdata selectData(String devBid) {
+        Map<String, String> params = new HashMap<>();
+        params.put("devBid", devBid);
+        return (IotRunHaoSfdata) mongoService.findOne(IotRunHaoSfdata.class, params, "sfdataCreatedDate", "desc");
+    }
+
+    @Override
+    public List<IotRunHaoSfdata> selectDataList(List<String> sfdataBidList) {
+        Map<String, Object> params = new HashMap<>();
+        params.put("newList_sfdataBid", sfdataBidList);
+        return mongoService.findAll(IotRunHaoSfdata.class, params);
+    }
+}

+ 213 - 0
src/main/java/com/yunfeiyun/agmp/iots/service/impl/IotSfElementfactorServiceImpl.java

@@ -0,0 +1,213 @@
+package com.yunfeiyun.agmp.iots.service.impl;
+
+import com.yunfeiyun.agmp.iot.common.domain.IotSfElementfactor;
+import com.yunfeiyun.agmp.iot.common.enums.EnumSfElementType;
+import com.yunfeiyun.agmp.iots.domain.IotSfElementfactorAlreadyListResVo;
+import com.yunfeiyun.agmp.iots.domain.IotSfElementfactorListReqVo;
+import com.yunfeiyun.agmp.iots.mapper.IotSfElementfactorMapper;
+import com.yunfeiyun.agmp.iots.service.IIotSfElementfactorService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.BeanUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+
+@Slf4j
+@Service
+public class IotSfElementfactorServiceImpl implements IIotSfElementfactorService {
+    @Autowired
+    private IotSfElementfactorMapper iotSfElementfactorMapper;
+    /**
+     * 查询水肥机要素列表
+     *
+     * @param  水肥机要素
+     * @return 水肥机要素集合
+     */
+    @Override
+    public List<IotSfElementfactor> selectIotSfElementfactorList(IotSfElementfactorListReqVo reqVo) {
+        return iotSfElementfactorMapper.selectIotSfElementfactorList(reqVo);
+    }
+
+
+    /**
+     * 批量插入水肥机要素
+     *
+     * @param factorList 水肥机要素列表
+     * @return 结果
+     */
+    @Transactional(rollbackFor = Exception.class)
+    @Override
+    public int batchInsertIotSfElementfactor(List<IotSfElementfactor> factorList) {
+        return iotSfElementfactorMapper.batchInsertIotSfElementfactor(factorList);
+    }
+
+    /**
+     * 删除水肥机要素
+     *
+     * @param sfBid 水肥机要素ID
+     * @return 结果
+     */
+    @Override
+    public int deleteIotSfElementfactorBySfBid(String sfBid) {
+        return iotSfElementfactorMapper.deleteIotSfElementfactorBySfBid(sfBid);
+    }
+
+    /**
+     * 批量删除水肥机要素
+     *
+     * @param sfBidList@return 结果
+     */
+
+    @Transactional(rollbackFor = Exception.class)
+    @Override
+    public int batchDeleteIotSfElementfactorBySfBidList(List<String> sfBidList) {
+        return iotSfElementfactorMapper.batchDeleteIotSfElementfactorBySfBidList(sfBidList);
+    }
+
+//    /**
+//     * 查询已配置要素列表
+//     *
+//     * @param reqVo
+//     * @return
+//     */
+//    @Override
+//    public List<IotSfElementfactorAlreadyListResVo> getAlreadyElementList(IotSfElementfactorListReqVo reqVo) {
+//        List<IotSfElementfactor> elementfactorList = selectIotSfElementfactorList(reqVo);
+//        Map<String, IotSfElementfactorAlreadyListResVo> eleMap = new LinkedHashMap<>();
+//        for (IotSfElementfactor elementfactor : elementfactorList) {
+//            String sfBid = elementfactor.getSfBid();
+//            String sfParentBid = elementfactor.getSfParentBid();
+//
+//            IotSfElementfactorAlreadyListResVo eleResVo = new IotSfElementfactorAlreadyListResVo();
+//            BeanUtils.copyProperties(elementfactor, eleResVo);
+//
+//            IotSfElementfactorAlreadyListResVo parentInfo = eleMap.get(sfParentBid);
+//            if(parentInfo == null){
+//                eleResVo.setChildrenList(new ArrayList<>());
+//                eleMap.put(sfBid, eleResVo);
+//            }else{
+//                parentInfo.getChildrenList().add(eleResVo);
+//            }
+//        }
+//        return new ArrayList<>(eleMap.values());
+//    }
+
+    /**
+     * 查询已配置要素列表
+     *
+     * @param factorList@return
+     */
+    @Override
+    public List<IotSfElementfactorAlreadyListResVo> getAlreadyElementList(List<IotSfElementfactor> factorList) {
+        Map<String, IotSfElementfactorAlreadyListResVo> eleMap = new LinkedHashMap<>();
+        for (IotSfElementfactor elementfactor : factorList) {
+            String sfBid = elementfactor.getSfBid();
+            String sfParentBid = elementfactor.getSfParentBid();
+
+            IotSfElementfactorAlreadyListResVo eleResVo = new IotSfElementfactorAlreadyListResVo();
+            BeanUtils.copyProperties(elementfactor, eleResVo);
+
+            IotSfElementfactorAlreadyListResVo parentInfo = eleMap.get(sfParentBid);
+            if(parentInfo == null){
+                eleResVo.setChildrenList(new ArrayList<>());
+                eleMap.put(sfBid, eleResVo);
+            }else{
+                parentInfo.getChildrenList().add(eleResVo);
+            }
+        }
+        return new ArrayList<>(eleMap.values());
+    }
+
+    /**
+     * 查询泵类已配置原始要素列表
+     *
+     * @param reqVo
+     * @return
+     */
+    @Override
+    public List<IotSfElementfactor> selectIotSfElementfactorListByPump(IotSfElementfactorListReqVo reqVo) {
+        List<String> sfTypeList = new ArrayList<>();
+        sfTypeList.add(EnumSfElementType.WATER_SOURCE.getCode());
+        sfTypeList.add(EnumSfElementType.FERTILIZER.getCode());
+        sfTypeList.add(EnumSfElementType.SUCTION.getCode());
+        sfTypeList.add(EnumSfElementType.MIXING.getCode());
+        sfTypeList.add(EnumSfElementType.FERTILIZER_BUCKET.getCode());
+        reqVo.setSfTypeList(sfTypeList);
+        return selectIotSfElementfactorList(reqVo);
+    }
+
+    /**
+     * 查询灌区已配置原始要素列表
+     *
+     * @param reqVo
+     * @return
+     */
+    @Override
+    public List<IotSfElementfactor> selectIotSfElementfactorListByGroup(IotSfElementfactorListReqVo reqVo) {
+        List<String> sfTypeList = new ArrayList<>();
+        sfTypeList.add(EnumSfElementType.GROUP.getCode());
+        reqVo.setSfTypeList(sfTypeList);
+        return selectIotSfElementfactorList(reqVo);
+    }
+
+    /**
+     * 查询电磁阀已配置原始要素列表
+     *
+     * @param reqVo
+     * @return
+     */
+    @Override
+    public List<IotSfElementfactor> selectIotSfElementfactorListByValve(IotSfElementfactorListReqVo reqVo) {
+        List<String> sfTypeList = new ArrayList<>();
+        sfTypeList.add(EnumSfElementType.SOLENOID_VALVE.getCode());
+        reqVo.setSfTypeList(sfTypeList);
+        return selectIotSfElementfactorList(reqVo);
+    }
+
+    /**
+     * 查询泵类已配置要素列表
+     *
+     * @param reqVo
+     * @return
+     */
+    @Override
+    public List<IotSfElementfactorAlreadyListResVo> getPumpAlreadyElementList(IotSfElementfactorListReqVo reqVo) {
+        List<IotSfElementfactor> elementfactorList = selectIotSfElementfactorListByPump(reqVo);
+        return getAlreadyElementList(elementfactorList);
+    }
+
+    /**
+     * 查询灌区已配置要素列表
+     *
+     * @param reqVo
+     * @return
+     */
+    @Override
+    public List<IotSfElementfactorAlreadyListResVo> getGroupAlreadyElementList(IotSfElementfactorListReqVo reqVo) {
+        List<String> sfTypeList = new ArrayList<>();
+        sfTypeList.add(EnumSfElementType.GROUP.getCode());
+        sfTypeList.add(EnumSfElementType.SOLENOID_VALVE.getCode());
+        reqVo.setSfTypeList(sfTypeList);
+
+        List<IotSfElementfactor> elementfactorList = selectIotSfElementfactorList(reqVo);
+        return getAlreadyElementList(elementfactorList);
+    }
+
+    /**
+     * 查询电磁阀已配置要素列表
+     *
+     * @param reqVo
+     * @return
+     */
+    @Override
+    public List<IotSfElementfactorAlreadyListResVo> getValveAlreadyElementList(IotSfElementfactorListReqVo reqVo) {
+        List<IotSfElementfactor> elementfactorList = selectIotSfElementfactorListByValve(reqVo);
+        return getAlreadyElementList(elementfactorList);
+    }
+}

+ 52 - 0
src/main/java/com/yunfeiyun/agmp/iots/service/impl/IotSfIrrigationRecordServiceImpl.java

@@ -0,0 +1,52 @@
+package com.yunfeiyun.agmp.iots.service.impl;
+
+import com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord;
+import com.yunfeiyun.agmp.iots.domain.IotSfIrrigationRecordListReqVo;
+import com.yunfeiyun.agmp.iots.mapper.IotSfIrrigationRecordMapper;
+import com.yunfeiyun.agmp.iots.service.IIotSfIrrigationRecordService;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import java.util.List;
+
+@Service
+public class IotSfIrrigationRecordServiceImpl implements IIotSfIrrigationRecordService {
+
+    @Autowired
+    private IotSfIrrigationRecordMapper irrigationRecordMapper;
+
+    @Override
+    public int insertIrrigationRecord(IotSfIrrigationRecord record) {
+        return irrigationRecordMapper.insertIrrigationRecord(record);
+    }
+
+    @Override
+    public int batchInsertIotSfIrrigationRecord(List<IotSfIrrigationRecord> recordList) {
+        return irrigationRecordMapper.batchInsertIotSfIrrigationRecord(recordList);
+    }
+
+    @Override
+    public IotSfIrrigationRecord selectIrrigationRecordByBid(String rcdBid) {
+        return irrigationRecordMapper.selectIrrigationRecordByBid(rcdBid);
+    }
+
+    @Override
+    public int updateIrrigationRecord(IotSfIrrigationRecord record) {
+        return irrigationRecordMapper.updateIrrigationRecord(record);
+    }
+
+    @Override
+    public int batchUpdateIrrigationRecord(List<IotSfIrrigationRecord> recordList) {
+        return irrigationRecordMapper.batchUpdateIrrigationRecord(recordList);
+    }
+
+    @Override
+    public void deleteIrrigationRecord(String rcdBid) {
+        return;
+    }
+
+    @Override
+    public List<IotSfIrrigationRecord> selectIrrigationRecordList(IotSfIrrigationRecordListReqVo record) {
+        return irrigationRecordMapper.selectIrrigationRecordList(record);
+    }
+}

+ 13 - 0
src/main/java/com/yunfeiyun/agmp/iots/startup/MongoStartup.java

@@ -163,6 +163,19 @@ public class MongoStartup {
         return mongodbIndexEntity;
     }
 
+    public MongodbIndexEntity iotIotRunHaoSfdataCreateIndex() {
+        log.info("开始创建IotRunHaoSfdata索引");
+        List<String[]> indexNameList = new ArrayList<>();
+        indexNameList.add(new String[]{"devBid"});
+        indexNameList.add(new String[]{"devBid", "sfdataCreatedDate"});
+
+        MongodbIndexEntity mongodbIndexEntity = new MongodbIndexEntity();
+        mongodbIndexEntity.setIotBaseEntity(IotRunHaoSfdata.class);
+        mongodbIndexEntity.setIndexNameList(indexNameList);
+
+        return mongodbIndexEntity;
+    }
+
     @PostConstruct
     public void start() {
         log.info("开始创建mongodb索引");

+ 3 - 2
src/main/java/com/yunfeiyun/agmp/iots/task/YbqScheduler.java

@@ -57,6 +57,7 @@ public class YbqScheduler {
     YbqTypeConst[] ybqTypeConsts = {
             YbqTypeConst.YBQ_XM_CMB,
             YbqTypeConst.YBQ_SD_DWB,
+            YbqTypeConst.YBQ_XM_TXB,
             YbqTypeConst.YBQ_YM_DBB,
     };
 
@@ -71,9 +72,9 @@ public class YbqScheduler {
      * 每天12点:0 0 8 * * ?
      * 每天8点:10分 0 10 8 * * ?
      * 对方:每天8点更新
-     * 我们:每天八点半更新
+     * 我们:每天8:10更新
      */
-    @Scheduled(cron = "0 30 8 * * ?")
+    @Scheduled(cron = "0 10 8 * * ?")
     public void synPredictedData() {
         String startDate = DateUtils.dateTime();
         log.info("【开始执行】同步预测数据任务,当前时间: {}", startDate);

+ 12 - 4
src/main/java/com/yunfeiyun/agmp/iots/warn/job/WarnJob.java

@@ -45,7 +45,7 @@ public class WarnJob {
      */
     @Scheduled(cron = "0 0/20 * * * ?")
     public void pestWarnJob20() {
-        if("0".equals(cbdDateDiff)){
+        if ("0".equals(cbdDateDiff)) {
             log.info("【设备预警】测报类检测,0 0/20 * * * ?  临时配置 每20分钟执行一次");
             // 处理虫害
             warnService.processWarningPestData();
@@ -55,16 +55,24 @@ public class WarnJob {
     }
 
     /**
-     * 每天八点
+     * 【虫害】每天八点
      */
     @Scheduled(cron = "0 0 8 * * ?")
     public void pestWarnJob08() {
-        log.info("【设备预警】测报类检测,0 0 8 * * ?  每天8点执行一次");
+        log.info("【设备预警】测报类检测,0 0 9 * * ?  每天9点执行一次");
         // 处理虫害
         warnService.processWarningPestData();
+
+    }
+
+    /**
+     * 【病害】每天八点::1出发病害
+     */
+    @Scheduled(cron = "0 15 8 * * ?")
+    public void pestWarnJob0815() {
+        log.info("【设备预警】测报类检测,0 15 8 * * ?  每天8:15 执行一次");
         // 处理病害
         warnService.processWarningDiseaseData();
 
     }
-
 }

+ 14 - 4
src/main/java/com/yunfeiyun/agmp/iots/warn/mapper/IotWarnBussinessMapper.java

@@ -1,9 +1,6 @@
 package com.yunfeiyun.agmp.iots.warn.mapper;
 
-import com.yunfeiyun.agmp.iot.common.domain.IotWarnconfig;
-import com.yunfeiyun.agmp.iot.common.domain.IotWarncount;
-import com.yunfeiyun.agmp.iot.common.domain.IotWarnindicator;
-import com.yunfeiyun.agmp.iot.common.domain.IotWarnlog;
+import com.yunfeiyun.agmp.iot.common.domain.*;
 import com.yunfeiyun.agmp.iots.warn.model.IotWarnconfigDevVo;
 import com.yunfeiyun.agmp.iots.warn.model.WarnConfigInfo;
 import org.apache.ibatis.annotations.Param;
@@ -73,6 +70,9 @@ public interface IotWarnBussinessMapper {
      */
     List<IotWarnindicator> selectCbdIndicatorAllList();
 
+    List<IotWarnindicator> selectYbqIndicatorAllList();
+    List<IotWarnconfigDevVo> selectIotWarnconfigYbqDevList();
+
 
     /**
      * 查询测报灯设备所有告警配置信息列表
@@ -80,4 +80,14 @@ public interface IotWarnBussinessMapper {
      * @return
      */
     List<IotWarnconfigDevVo> selectIotWarnconfigCbdDevList();
+
+    IotDevice selectDeviceById(@Param("devBid") String devBid);
+
+    IotWarnpolicy selectWarnPolicy(@Param("configId") String configId);
+
+    List<IotWarnreceiver> selectWarnReceiverByConfigId(String configId);
+
+    void updateWarnLogSendStatus(@Param("time") String time, @Param("wlBid") String wlBid, @Param("status") String status);
+
+    IotWarnlog getLastedUnSendWarnLog(String wcBid);
 }

+ 1 - 1
src/main/java/com/yunfeiyun/agmp/iots/warn/model/IotWarnconfigCbdInfoVo.java

@@ -8,7 +8,7 @@ import lombok.Data;
 import java.util.List;
 
 @Data
-public class IotWarnconfigCbdInfoVo extends IotWarnconfig {
+public class IotWarnconfigInfoVo extends IotWarnconfig {
 
     private IotDevice iotDevice;
 

+ 10 - 0
src/main/java/com/yunfeiyun/agmp/iots/warn/model/WarnResult.java

@@ -2,8 +2,11 @@ package com.yunfeiyun.agmp.iots.warn.model;
 
 import com.yunfeiyun.agmp.iot.common.domain.IotBaseEntity;
 import com.yunfeiyun.agmp.iot.common.domain.IotWarnconfig;
+import com.yunfeiyun.agmp.iot.common.domain.IotWarnpolicy;
 import lombok.Data;
 
+import java.util.List;
+
 /**
  * 预警结果,不是对应数据库层面的
  */
@@ -53,6 +56,13 @@ public class WarnResult extends IotBaseEntity {
      * 是不是离线指标
      */
     private boolean isOffline;
+    private IotWarnpolicy iotWarnpolicy;
+
+    private List<String> receiverIds;
+    /**
+     * 预警日志id
+     */
+    private String wlBid;
 
     public WarnResult(boolean isTriggered, String message) {
         this.isTriggered = isTriggered;

+ 64 - 6
src/main/java/com/yunfeiyun/agmp/iots/warn/service/IotWarnBussinessService.java

@@ -101,11 +101,11 @@ public class IotWarnBussinessService {
         return result;
     }
 
-    public int insertWarnRecord(IotWarnlog iotWarnlog) {
+    public String insertWarnRecord(IotWarnlog iotWarnlog) {
         log.info("插入预警记录: {}", iotWarnlog);
         int result = iotWarnBussinessMapper.insertWarnRecord(iotWarnlog);
         log.info("插入预警记录结果: {}", result);
-        return result;
+        return iotWarnlog.getWlBid();
     }
 
     public List<WarnConfigInfo> selectIotWarnConfigInfoList(WarnConfigInfo warnConfigInfo) {
@@ -330,10 +330,10 @@ public class IotWarnBussinessService {
      * @return
      */
 
-    public List<IotWarnconfigCbdInfoVo> selectIotWarnconfigCbdInfoList() {
+    public List<IotWarnconfigInfoVo> selectIotWarnconfigCbdInfoList() {
         List<IotWarnconfigDevVo> iotWarnconfigDevVoList = iotWarnBussinessMapper.selectIotWarnconfigCbdDevList();
         Map<String, List<IotWarnindicator>> iotWarnindicatorMap = selectCbdIndicatorAllMap();
-        List<IotWarnconfigCbdInfoVo> iotWarnconfigCbdInfoVoList = new ArrayList<>();
+        List<IotWarnconfigInfoVo> iotWarnconfigCbdInfoVoList = new ArrayList<>();
         List<String> parentbidList = new ArrayList<>();
         for (IotWarnconfigDevVo iotWarnconfigDevVo : iotWarnconfigDevVoList) {
             String wcBid = iotWarnconfigDevVo.getWcBid();
@@ -347,7 +347,7 @@ public class IotWarnBussinessService {
             IotDevice iotDevice = new IotDevice();
             BeanUtils.copyProperties(iotWarnconfigDevVo, iotDevice);
 
-            IotWarnconfigCbdInfoVo iotWarnconfigCbdInfoVo = new IotWarnconfigCbdInfoVo();
+            IotWarnconfigInfoVo iotWarnconfigCbdInfoVo = new IotWarnconfigInfoVo();
             BeanUtils.copyProperties(iotWarnconfigDevVo, iotWarnconfigCbdInfoVo);
             iotWarnconfigCbdInfoVo.setIotDevice(iotDevice);
             iotWarnconfigCbdInfoVo.setIotWarnindicatorList(iotWarnindicatorList);
@@ -357,7 +357,7 @@ public class IotWarnBussinessService {
             IotWarnindicator selectIotWarnindicator = new IotWarnindicator();
             selectIotWarnindicator.setWiParentbidList(parentbidList);
             Map<String, List<IotWarnindicator>> iotMap = selectIotWarnindicatorMap(selectIotWarnindicator);
-            for (IotWarnconfigCbdInfoVo cbdInfoVo : iotWarnconfigCbdInfoVoList) {
+            for (IotWarnconfigInfoVo cbdInfoVo : iotWarnconfigCbdInfoVoList) {
                 String wcBid = cbdInfoVo.getWcBid();
                 List<IotWarnindicator> pestDetailList = iotMap.get(wcBid);
                 List<IotWarnindicator> iotWarnindicatorList = cbdInfoVo.getIotWarnindicatorList();
@@ -371,4 +371,62 @@ public class IotWarnBussinessService {
         }
         return iotWarnconfigCbdInfoVoList;
     }
+
+    /**
+     * 将Ybq设备告警的配置信息。设备、要素合并一起
+     *
+     * @return
+     */
+    public List<IotWarnconfigInfoVo> selectIotWarnconfigYbqInfoList() {
+        List<IotWarnconfigDevVo> iotWarnconfigDevVoList = iotWarnBussinessMapper.selectIotWarnconfigYbqDevList();
+        Map<String, List<IotWarnindicator>> iotWarnindicatorMap = selectYbqIndicatorAllMap();
+        List<IotWarnconfigInfoVo> iotWarnconfigCbdInfoVoList = new ArrayList<>();
+        for (IotWarnconfigDevVo iotWarnconfigDevVo : iotWarnconfigDevVoList) {
+            String wcBid = iotWarnconfigDevVo.getWcBid();
+            List<IotWarnindicator> iotWarnindicatorList = iotWarnindicatorMap.get(wcBid);
+            IotDevice iotDevice = new IotDevice();
+            BeanUtils.copyProperties(iotWarnconfigDevVo, iotDevice);
+
+            IotWarnconfigInfoVo iotWarnconfigCbdInfoVo = new IotWarnconfigInfoVo();
+            BeanUtils.copyProperties(iotWarnconfigDevVo, iotWarnconfigCbdInfoVo);
+            iotWarnconfigCbdInfoVo.setIotDevice(iotDevice);
+            iotWarnconfigCbdInfoVo.setIotWarnindicatorList(iotWarnindicatorList);
+            iotWarnconfigCbdInfoVoList.add(iotWarnconfigCbdInfoVo);
+        }
+        return iotWarnconfigCbdInfoVoList;
+    }
+
+    public IotDevice selectDeviceById(String devBid) {
+        return iotWarnBussinessMapper.selectDeviceById(devBid);
+    }
+
+    public Map<String, List<IotWarnindicator>> selectYbqIndicatorAllMap() {
+        List<IotWarnindicator> iotWarnindicators = iotWarnBussinessMapper.selectYbqIndicatorAllList();
+        Map<String, List<IotWarnindicator>> iotWarnindicatorMap = new LinkedHashMap<>();
+        for (IotWarnindicator iotWarnindicator : iotWarnindicators) {
+            String wcBid = iotWarnindicator.getWcBid();
+            if (!iotWarnindicatorMap.containsKey(wcBid)) {
+                iotWarnindicatorMap.put(wcBid, new ArrayList<>());
+            }
+            iotWarnindicatorMap.get(wcBid).add(iotWarnindicator);
+        }
+        return iotWarnindicatorMap;
+    }
+
+    public IotWarnpolicy selectWarnPolicy(String configId) {
+        return iotWarnBussinessMapper.selectWarnPolicy(configId);
+    }
+
+    public List<IotWarnreceiver> selectWarnReceiverByConfigId(String configId) {
+
+        return iotWarnBussinessMapper.selectWarnReceiverByConfigId(configId);
+    }
+
+    public void updateWarnLogSendStatus(String time, String wlBid, String status) {
+        iotWarnBussinessMapper.updateWarnLogSendStatus(time, wlBid, status);
+    }
+
+    public IotWarnlog getLastedUnSendWarnLog(String wcBid) {
+        return iotWarnBussinessMapper.getLastedUnSendWarnLog(wcBid);
+    }
 }

+ 235 - 0
src/main/java/com/yunfeiyun/agmp/iots/warn/service/MsgService.java

@@ -0,0 +1,235 @@
+package com.yunfeiyun.agmp.iots.warn.service;
+
+import com.alibaba.fastjson2.JSONObject;
+import com.yunfeiyun.agmp.common.enums.warn.MsgBusType;
+import com.yunfeiyun.agmp.common.enums.warn.MsgType;
+import com.yunfeiyun.agmp.common.enums.warn.WarnLogSendStatus;
+import com.yunfeiyun.agmp.common.framework.message.MessageDto;
+import com.yunfeiyun.agmp.common.framework.mq.rabbitmq.enums.AgmpActionEnums;
+import com.yunfeiyun.agmp.common.framework.mq.rabbitmq.model.SynAgmpInfoDto;
+import com.yunfeiyun.agmp.common.utils.DateUtils;
+import com.yunfeiyun.agmp.common.utils.StringUtils;
+import com.yunfeiyun.agmp.iot.common.constant.devicetype.IotDeviceDictEnum;
+import com.yunfeiyun.agmp.iot.common.constant.devicetype.IotDeviceTypeLv1Enum;
+import com.yunfeiyun.agmp.iot.common.domain.IotWarnconfig;
+import com.yunfeiyun.agmp.iot.common.domain.IotWarnlog;
+import com.yunfeiyun.agmp.iot.common.domain.IotWarnpolicy;
+import com.yunfeiyun.agmp.iot.common.domain.IotWarnreceiver;
+import com.yunfeiyun.agmp.iots.mq.provider.AgmpIotMqProviderService;
+import com.yunfeiyun.agmp.iots.mq.provider.AgmpMqProviderService;
+import com.yunfeiyun.agmp.iots.warn.model.WarnResult;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.List;
+import java.util.stream.Collectors;
+
+/**
+ * 消息中心
+ * 重复次数服务产生告警信息后发给消息中心发送消息
+ */
+@Service
+@Slf4j
+public class MsgService {
+    @Autowired
+    private IotWarnBussinessService iotWarnBussinessService;
+
+    @Autowired
+    private AgmpIotMqProviderService agmpMqProviderService;
+
+    /**
+     * 1. 判断是否触发
+     * 2. 如果已经触发预警,获取配置信息IotWarnpolicy,检查是即使发送还是周期性发送wpType:0即时推送,1选定时间
+     * 3. 如果是及时发送
+     * 3.1 如果是站内信web,直接入库,代表发送
+     * 3.2 首先入库,状态待发送,然后放入队列,异步发送
+     * 4. 如果是选定时间
+     * 4.1 获取需要发送的时间段,还要获取上一次发送的时间,以及配置信息的发送周期
+     * 4.2 如果在这个发送时间段,离上次发送消息已经过了一个周期,则发送,否则不发送
+     *
+     * @param warnResult
+     */
+    public void handleWarn(WarnResult warnResult) {
+
+        try {
+            if (!warnResult.isTriggered()) {
+                return;
+            }
+            String messageId = warnResult.getMessageId();
+            IotWarnpolicy iotWarnpolicy = iotWarnBussinessService.selectWarnPolicy(warnResult.getConfigId());
+            if (iotWarnpolicy == null) {
+                log.info("未找到配置信息,配置id:{}", warnResult.getConfigId());
+                return;
+            }
+            warnResult.setIotWarnpolicy(iotWarnpolicy);
+            List<IotWarnreceiver> iotWarnreceiverList = iotWarnBussinessService.selectWarnReceiverByConfigId(warnResult.getConfigId());
+            if (iotWarnreceiverList.isEmpty()) {
+                iotWarnreceiverList = new ArrayList<>();
+            }
+            List<String> userIds = iotWarnreceiverList.stream().map(IotWarnreceiver::getWruserId).collect(Collectors.toList());
+            warnResult.setReceiverIds(userIds);
+            String wpType = iotWarnpolicy.getWpType();
+            boolean reSend = false;
+            if (wpType.equals("0")) {
+                // 即时推送
+                log.info("【告警通知】消息标识{}:当前设备ID:{},即时推送", messageId, warnResult.getDevId());
+                //如果设置有推送频率,结合最后一次发送时间,如果上次距离当下超过发送频率,那就可以发送
+                // 发送频率,每N时,数字类型
+                String wpFrequency = StringUtils.isEmpty(iotWarnpolicy.getWpFrequency())?"0":iotWarnpolicy.getWpFrequency();
+                // 判断是否到了发送时间
+                IotWarnlog lastedUnSendWarnLog = getLastedUnSendWarnLog(warnResult.getConfigId());
+                if (lastedUnSendWarnLog != null) {
+                    // 根据最后一次发送时间,结合发送频率,如果上次距离当下超过发送频率,那就可以发送
+                    if (!DateUtils.isBetween(new Date(), DateUtils.parseDate(lastedUnSendWarnLog.getWlCreateddate()), DateUtils.addHours(DateUtils.parseDate(lastedUnSendWarnLog.getWlCreateddate()), Integer.parseInt(wpFrequency)))) {
+                        saveWarnMsg(warnResult);
+                    } else {
+                        reSend = true;
+                        log.info("还没到发送时间,下次发送时间:{}", DateUtils.addHours(DateUtils.parseDate(lastedUnSendWarnLog.getWlCreateddate()), Integer.parseInt(wpFrequency)).toString());
+                    }
+                } else {
+                    saveWarnMsg(warnResult);
+                }
+            } else if (wpType.equals("1")) {
+                // 选定时间
+                String deliveryTimePeriod = iotWarnpolicy.getDeliveryTimePeriod();
+                String startTime = deliveryTimePeriod.split("-")[0];
+                String endTime = deliveryTimePeriod.split("-")[1];
+                Date startDate = DateUtils.parseDate(DateUtils.getDate()+" "+startTime);
+                Date endDate = DateUtils.parseDate(DateUtils.getDate()+" "+endTime);
+                log.info("【告警通知】消息标识{}:当前设备ID:{},时间段推送:deliveryTimePeriod{},start:{},end:{}", messageId, warnResult.getDevId(),deliveryTimePeriod,startDate,endDate);
+                if (DateUtils.isBetween(new Date(), startDate, endDate)) {
+                    // 发送频率,每N时,数字类型
+                    String wpFrequency = StringUtils.isEmpty(iotWarnpolicy.getWpFrequency())?"0":iotWarnpolicy.getWpFrequency();
+                    // 判断是否到了发送时间
+                    IotWarnlog lastedUnSendWarnLog = getLastedUnSendWarnLog(warnResult.getConfigId());
+                    if (lastedUnSendWarnLog != null) {
+                        // 根据最后一次发送时间,结合发送频率,如果上次距离当下超过发送频率,那就可以发送
+                        if (!DateUtils.isBetween(new Date(), DateUtils.parseDate(lastedUnSendWarnLog.getWlCreateddate()), DateUtils.addHours(DateUtils.parseDate(lastedUnSendWarnLog.getWlCreateddate()), Integer.parseInt(wpFrequency)))) {
+                            saveWarnMsg(warnResult);
+                        } else {
+                            reSend = true;
+                            log.info("还没到发送时间,下次发送时间:{}", DateUtils.addHours(DateUtils.parseDate(lastedUnSendWarnLog.getWlCreateddate()), Integer.parseInt(wpFrequency)).toString());
+                        }
+                    } else {
+                        saveWarnMsg(warnResult);
+                    }
+                } else {
+                    reSend = true;
+                }
+            }
+            //统一处理补发逻辑
+            reSendWarn(warnResult.getWlBid(), reSend, warnResult.getConfigId(), iotWarnpolicy.getDeliveryTimePeriod(), iotWarnpolicy.getWpFrequency());
+        } catch (Exception e) {
+            log.error("handleWarn error", e);
+        }
+    }
+
+    /**
+     * 补发逻辑
+     *
+     * @param wlId
+     * @param reSendStatus
+     * @param wlBid
+     * @param deliveryTimePeriod
+     * @param wpFrequency
+     */
+    void reSendWarn(String wlId, boolean reSendStatus, String wlBid, String deliveryTimePeriod, String wpFrequency) {
+        if (reSendStatus) {
+            //标记需要补发
+        } else {
+            //标记不需要补发
+        }
+    }
+
+    void saveWarnMsg(WarnResult warnResult) {
+        SynAgmpInfoDto synAgmpInfoDto = new SynAgmpInfoDto();
+        synAgmpInfoDto.setAction(AgmpActionEnums.AGMP_SEND_MSG.getCode());
+        synAgmpInfoDto.setData(JSONObject.from(resolvePtsMsgByIotWarnlog(warnResult)));
+        synAgmpInfoDto.setDesc("发送消息");
+        //标记告警记录,已发送;
+        iotWarnBussinessService.updateWarnLogSendStatus(DateUtils.dateTimeNow(), warnResult.getWlBid(), WarnLogSendStatus.SEND_SUCCESS.getCode());
+        agmpMqProviderService.sendToAgmpIot(JSONObject.toJSONString(synAgmpInfoDto));
+    }
+
+    public IotWarnlog getLastedUnSendWarnLog(String wcBid) {
+        // 获取最新的一条未发送的告警
+        IotWarnlog iotWarnlog = iotWarnBussinessService.getLastedUnSendWarnLog(wcBid);
+        return iotWarnlog;
+    }
+
+
+    /**
+     * 告警日志转为消息对象
+     */
+    public MessageDto resolvePtsMsgByIotWarnlog(WarnResult warnResult) {
+        MessageDto messageDto = new MessageDto();
+        IotWarnconfig iotWarnconfig = warnResult.getConfig();
+        IotWarnpolicy iotWarnpolicy = warnResult.getIotWarnpolicy();
+        // 直接映射或推测可能的对应关系
+        messageDto.setMsgbatchId(warnResult.getMessageId());  // 假设消息标识作为消息批次标识
+        messageDto.setMsgbatchContent(warnResult.getMessage());  // 使用预警消息作为消息批次内容
+        messageDto.setMsgbatchContenttype("1");  // 假定使用固定内容,具体依据业务需求确定
+        messageDto.setMsgbatchMsgtype(MsgType.WARN_INFO.getCode());
+        messageDto.setMsgbatchBiztype(resolveMsgBusTypeByDevtype(warnResult.getDevtypeBid()).getCode());  // 业务类型,这里假设为预警配置ID
+        messageDto.setMsgbatchBizobj(warnResult.getDevId());  // 设备ID作为业务对象
+        messageDto.setMsgbatchLevel(iotWarnconfig.getWcLevel());  // 消息等级,这里默认为1,实际应用中应依据具体情况设置
+        messageDto.setMsgbatchSource("iots");  // 消息来源,假设来自物联网系统
+        messageDto.setMsgbatchSender("system");  // 发送者,这里假设为系统发送
+        messageDto.setMsgbatchChannel(iotWarnpolicy.getWpChannel());  // 通知渠道,这里假设通过电子邮件发送,需要根据实际情况调整
+        JSONObject extra = new JSONObject();
+        extra.put("location", "/iotm/warning/record");
+        extra.put("permission", "iotm:warning:record:list");
+        extra.put("handleId", warnResult.getWlBid());
+        messageDto.setMsgbatchExtra(extra.toJSONString());
+        messageDto.setMsgbatchHandler(warnResult.getReceiverIds());  // 处理人暂时为空,根据业务逻辑补充
+        messageDto.setMsgbatchHandlerType("1");  // 处理人暂时为空,根据业务逻辑补充
+        messageDto.setMsgbatchCreateddate(DateUtils.dateTimeNow());  // 创建时间
+        messageDto.setMsgbatchTitle("设备告警");  // 消息标题,这里假设为物联网设备警告
+        messageDto.setMsgbatchTarget("PTS");  // 目标子系统,这里假设为目标设备型号ID
+        messageDto.setWlBid(warnResult.getWlBid());  // 物联网预警ID
+        return messageDto;
+    }
+
+    //根据设备型号判断对应那种告警类型
+    public MsgBusType resolveMsgBusTypeByDevtype(String devtypeBid) {
+        String devClass = IotDeviceDictEnum.getLv1CodeByCode(devtypeBid);
+        IotDeviceTypeLv1Enum iotDeviceTypeLv1Enum = IotDeviceTypeLv1Enum.findEnumByCode(devClass);
+        switch (iotDeviceTypeLv1Enum) {
+            case QXZ: {
+                return MsgBusType.WARN_QX;
+            }
+            case SQZ:
+            case GSSQ: {
+                return MsgBusType.WARN_SQ;
+            }
+            //病害
+            case YBQ_DWB:
+            case YBQ_CMB:
+            case YBQ_DBB:
+            case YBQ_TXB:
+            case YBQ_BFB: {
+                return MsgBusType.WARN_DISEASE;
+            }
+            //虫害
+            case CBD:
+                return MsgBusType.WARN_PEST;
+            default:
+                return MsgBusType.WARN_DEVICE_COMMON;
+        }
+    }
+
+
+    /**
+     * 插入消息内容
+     */
+    public void saveWarnMsg(IotWarnlog iotWarnlog) {
+
+    }
+
+/**
+ * 定时任务:获取到了需要发送的时间段,获取最新的一条未发送的告警进行发送
+ */
+}

+ 12 - 2
src/main/java/com/yunfeiyun/agmp/iots/warn/service/ReCountService.java

@@ -1,6 +1,7 @@
 package com.yunfeiyun.agmp.iots.warn.service;
 
 import com.yunfeiyun.agmp.common.constant.ErrorCode;
+import com.yunfeiyun.agmp.common.enums.warn.WarnLogSendStatus;
 import com.yunfeiyun.agmp.common.exception.BizException;
 import com.yunfeiyun.agmp.common.utils.DateUtils;
 import com.yunfeiyun.agmp.common.utils.StringUtils;
@@ -20,6 +21,8 @@ import org.springframework.stereotype.Service;
 public class ReCountService {
     @Autowired
     private IotWarnBussinessService iotWarnBussinessService;
+    @Autowired
+    private MsgService msgService;
 
     /**
      * 针对预警解析出来的结果进行处理
@@ -35,7 +38,10 @@ public class ReCountService {
         // 离线的只用直接入库即可
         if (warnResult.isTriggered() && warnResult.isOffline()) {
             log.info("【设备预警】消息标识{}:当前设备ID:{},离线告警", messageId, warnResult.getDevId(), warnResult.getConfigId(), warnResult);
-            iotWarnBussinessService.insertWarnRecord(buildWarnMessage(messageId, warnResult));
+            String wlBid = iotWarnBussinessService.insertWarnRecord(buildWarnMessage(messageId, warnResult));
+            //发给消息服务
+            warnResult.setWlBid(wlBid);
+            msgService.handleWarn(warnResult);
             return;
         }
         // 触发了预警,进行校验处理
@@ -46,7 +52,10 @@ public class ReCountService {
             log.info("【设备预警】消息标识{}:当前设备ID:{}, 预警配置ID:{} 的重复次数为: {}", messageId, warnResult.getDevId(), warnResult.getConfigId(), thisReCount);
             if (thisReCount + 1 > targetReCount) {
                 log.info("【设备预警】消息标识{}:达到或超过阈值, 不再增加重复次数,直接生成预警记录", messageId);
-                iotWarnBussinessService.insertWarnRecord(buildWarnMessage(messageId, warnResult));
+                String wlBid = iotWarnBussinessService.insertWarnRecord(buildWarnMessage(messageId, warnResult));
+                //发给消息服务
+                warnResult.setWlBid(wlBid);
+                msgService.handleWarn(warnResult);
             } else {
                 log.info("【设备预警】消息标识{}:未达阈值, 增加重复次数", messageId);
                 iotWarnBussinessService.incrementReCount(reCount, warnResult);
@@ -115,6 +124,7 @@ public class ReCountService {
         iotWarnlog.setWlData(warnResult.getReportData());
         iotWarnlog.setTid(config.getTid());
         iotWarnlog.setWcBid(config.getWcBid());
+        iotWarnlog.setWlSendmsgstatus(WarnLogSendStatus.SEND_WAIT.getCode());
         log.info("【设备预警】消息标识{}:构建预警信息: {}", messageId, iotWarnlog);
         return iotWarnlog;
     }

+ 165 - 0
src/main/java/com/yunfeiyun/agmp/iots/warn/service/WarnDiseaseService.java

@@ -0,0 +1,165 @@
+package com.yunfeiyun.agmp.iots.warn.service;
+
+import com.alibaba.fastjson2.JSONObject;
+import com.yunfeiyun.agmp.common.constant.ErrorCode;
+import com.yunfeiyun.agmp.common.exception.BizException;
+import com.yunfeiyun.agmp.common.utils.StringUtils;
+import com.yunfeiyun.agmp.iot.common.constant.devicetype.IotDeviceDictEnum;
+import com.yunfeiyun.agmp.iot.common.domain.IotDevice;
+import com.yunfeiyun.agmp.iot.common.domain.IotWarnconfig;
+import com.yunfeiyun.agmp.iot.common.domain.IotWarnindicator;
+import com.yunfeiyun.agmp.iot.common.enums.EnumWarnRuleOp;
+import com.yunfeiyun.agmp.iots.warn.model.*;
+import com.yunfeiyun.agmp.iots.warn.util.CompareUtil;
+import com.yunfeiyun.agmp.iots.warn.util.WarnMessageBuilderUtil;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.BeanUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import java.util.*;
+
+@Service
+@Slf4j
+public class WarnDiseaseService {
+    @Autowired
+    private IotWarnBussinessService iotWarnBussinessService;
+    @Autowired
+    private ReCountService reCountService;
+
+    /**
+     * 虫害处理入口方法,负责遍历所有设备配置信息并进行相应的处理。
+     */
+    public void process() {
+        List<IotWarnconfigInfoVo> iotWarnconfigInfoVos = iotWarnBussinessService.selectIotWarnconfigYbqInfoList();
+        for (IotWarnconfigInfoVo iotWarnconfigCbdInfoVo : iotWarnconfigInfoVos) {
+            try {
+                handleDevice(iotWarnconfigCbdInfoVo);
+            } catch (Exception e) {
+                e.printStackTrace();
+                log.error("处理设备 {} 时发生错误: {}", iotWarnconfigCbdInfoVo.getIotDevice().getDevBid(), e);
+            }
+        }
+    }
+
+    /**
+     * 处理单个设备的所有预警配置。
+     *
+     * @param iotWarnconfigCbdInfoVo 设备配置信息对象
+     */
+    private void handleDevice(IotWarnconfigInfoVo iotWarnconfigCbdInfoVo) {
+        IotDevice device = iotWarnconfigCbdInfoVo.getIotDevice();
+        String devBid = device.getDevBid();
+        List<IotWarnindicator> warnIndicators = iotWarnconfigCbdInfoVo.getIotWarnindicatorList();
+        IotWarnconfig warnConfig = new IotWarnconfig();
+        BeanUtils.copyProperties(iotWarnconfigCbdInfoVo, warnConfig);
+
+        log.info("开始处理设备ID为 {} 的预警配置", devBid);
+        for (IotWarnindicator indicator : warnIndicators) {
+            if ("computeValue".equals(indicator.getWiCode())) {
+                handleComputeValue(device, indicator, warnConfig);
+            }
+        }
+        log.info("完成设备ID为 {} 的所有预警配置处理", devBid);
+    }
+
+    /**
+     * 处理特定于“computeValue”的逻辑,并调用告警日志处理器。
+     *
+     * @param device     设备对象
+     * @param indicator  预警指标对象
+     * @param warnConfig 预警配置对象
+     */
+    private void handleComputeValue(IotDevice device, IotWarnindicator indicator, IotWarnconfig warnConfig) {
+        IotDevice detailedDevice = iotWarnBussinessService.selectDeviceById(device.getDevBid());
+        if (StringUtils.isNotEmpty(detailedDevice.getExtInfo())) {
+            JSONObject extInfo = JSONObject.parseObject(detailedDevice.getExtInfo());
+            Double computeValue = Double.parseDouble(extInfo.getOrDefault("computeValue", "0.0").toString());
+            WarnStatusDto warnStatusDto = buildWarnStatusDto(computeValue, indicator);
+            if (warnStatusDto != null) {
+                handlerAllWarnLog(device.getDevBid(), warnConfig, warnStatusDto, detailedDevice.getExtInfo());
+            } else {
+                log.warn("【设备预警】病虫害:设备id 不可为空. 设备ID: {}, 设备Code:{} , 未触发: ", detailedDevice.getDevBid(), detailedDevice.getDevCode());
+            }
+        } else {
+            log.warn("【设备预警】病虫害:设备id 不可为空. 设备ID: {}, 设备Code:{} , 最新数据为空: ", detailedDevice.getDevBid(), detailedDevice.getDevCode());
+        }
+
+    }
+
+    /**
+     * 根据当前值和预警指标构建WarnStatusDto对象。
+     *
+     * @param currentValue 当前值
+     * @param indicator    预警指标对象
+     * @return 构建好的WarnStatusDto对象
+     */
+    private WarnStatusDto buildWarnStatusDto(Double currentValue, IotWarnindicator indicator) {
+        if (!"0".equals(indicator.getWiStatus())) {
+            return null;
+        }
+        EnumWarnRuleOp op = EnumWarnRuleOp.findEnumByCode(indicator.getWiExpression());
+        boolean comparisonResult = CompareUtil.comp(currentValue.toString(), indicator.getWiExpression(), indicator.getWiValue());
+        WarnStatusDto dto = new WarnStatusDto();
+        dto.setDevType(IotDeviceDictEnum.getLv1NameByCode(indicator.getDevtypeBid()));
+        dto.setDevCode(indicator.getDevCode());
+        dto.setName(indicator.getWiName());
+        dto.setValue(currentValue.toString());
+        dto.setUnit(indicator.getWiUnit());
+        dto.setOpt(op.getName());
+        dto.setIndicatorValue(indicator.getWiValue());
+        dto.setWarn(comparisonResult);
+        log.debug("构建告警状态DTO: {}", dto);
+        return dto;
+    }
+
+    /**
+     * 处理所有的告警日志,包括生成警告消息、准备发送的消息参数,并最终发送告警。
+     *
+     * @param devBid 设备ID
+     * @param config 配置对象
+     * @param dto    告警状态数据传输对象
+     * @param data   设备额外信息
+     */
+    private void handlerAllWarnLog(String devBid, IotWarnconfig config, WarnStatusDto dto, String data) {
+        if (StringUtils.isEmpty(devBid)) {
+            log.error("【设备预警】病虫害:设备id 不可为空. 设备ID: {}, 配置ID: {}", devBid, config.getWcBid());
+            throw new BizException(ErrorCode.FAILURE.getCode(), "病虫害:设备id 不可为空");
+        }
+        log.info("开始处理设备ID为 {}, 配置ID为 {} 的预警信息", devBid, config.getWcBid());
+        String message = WarnMessageBuilderUtil.buildWarningMessage(dto.getDevType(), dto.getDevCode(),
+                dto.getName(), dto.getValue(), dto.getUnit(), dto.getOpt(), dto.getIndicatorValue());
+        log.info("根据类型生成警告消息: {}. 设备ID: {}, 配置ID: {}", message, devBid, config.getWcBid());
+
+        WarnResult result = createWarnResult(config, devBid, dto, data);
+        log.info("准备发送告警信息: {}. 设备ID: {}, 配置ID: {}", result.getMessage(), devBid, config.getWcBid());
+        reCountService.handlerMessage(result);
+        log.info("告警信息已发送完成. 设备ID: {}, 配置ID: {}", devBid, config.getWcBid());
+    }
+
+    /**
+     * 创建并初始化WarnResult对象,用于后续的告警消息发送。
+     *
+     * @param config 配置对象
+     * @param devBid 设备ID
+     * @param dto    告警状态数据传输对象
+     * @param data   设备额外信息
+     * @return 初始化后的WarnResult对象
+     */
+    private WarnResult createWarnResult(IotWarnconfig config, String devBid, WarnStatusDto dto, String data) {
+        WarnResult result = new WarnResult();
+        result.setMessageId(UUID.randomUUID().toString()); // 假设UUID方法用于生成唯一ID
+        result.setDevId(devBid);
+        result.setTid(config.getTid());
+        result.setConfigId(config.getWcBid());
+        result.setReportData(data);
+        result.setTargetReCount(config.getWcRepeatnum() == null ? 0 : config.getWcRepeatnum());
+        result.setDevtypeBid(config.getDevtypeBid());
+        result.setConfig(config);
+        result.setTriggered(dto.isWarn());
+        result.setMessage(WarnMessageBuilderUtil.buildWarningMessage(
+                dto.getDevType(), dto.getDevCode(), dto.getName(), dto.getValue(), dto.getUnit(), dto.getOpt(), dto.getIndicatorValue()
+        ));
+        return result;
+    }
+}

+ 3 - 3
src/main/java/com/yunfeiyun/agmp/iots/warn/service/WarnPestService.java

@@ -11,7 +11,7 @@ import com.yunfeiyun.agmp.iot.common.domain.IotWarnconfig;
 import com.yunfeiyun.agmp.iot.common.domain.IotWarnindicator;
 import com.yunfeiyun.agmp.iot.common.enums.EnumWarnRuleOp;
 import com.yunfeiyun.agmp.iots.service.IIotPestService;
-import com.yunfeiyun.agmp.iots.warn.model.IotWarnconfigCbdInfoVo;
+import com.yunfeiyun.agmp.iots.warn.model.IotWarnconfigInfoVo;
 import com.yunfeiyun.agmp.iots.warn.model.WarnPestDetailStatDto;
 import com.yunfeiyun.agmp.iots.warn.model.WarnResult;
 import com.yunfeiyun.agmp.iots.warn.model.WarnStatusDto;
@@ -45,8 +45,8 @@ public class WarnPestService {
      */
     public void process() {
         //单个配置单个设备的处理逻辑,外层需要遍历
-        List<IotWarnconfigCbdInfoVo> iotWarnconfigCbdInfoVoList = iotWarnBussinessService.selectIotWarnconfigCbdInfoList();
-        for (IotWarnconfigCbdInfoVo iotWarnconfigCbdInfoVo : iotWarnconfigCbdInfoVoList) {
+        List<IotWarnconfigInfoVo> iotWarnconfigCbdInfoVoList = iotWarnBussinessService.selectIotWarnconfigCbdInfoList();
+        for (IotWarnconfigInfoVo iotWarnconfigCbdInfoVo : iotWarnconfigCbdInfoVoList) {
             try {
                 IotDevice iotDevice = iotWarnconfigCbdInfoVo.getIotDevice();
                 String devBid = iotDevice.getDevBid();

+ 3 - 1
src/main/java/com/yunfeiyun/agmp/iots/warn/service/WarnService.java

@@ -50,6 +50,8 @@ public class WarnService {
     private IIotDevicefactorService iotDevicefactorService;
     @Autowired
     private WarnPestService warnPestService;
+    @Autowired
+    private WarnDiseaseService warnDiseaseService;
 
 
     /**
@@ -581,6 +583,6 @@ public class WarnService {
      * 病害
      */
     public void processWarningDiseaseData() {
-
+        warnDiseaseService.process();
     }
 }

+ 13 - 2
src/main/resources/application-dev.yml

@@ -52,6 +52,9 @@ user:
 
 # Spring配置
 spring:
+  autoconfigure:
+    exclude:
+      - org.springframework.boot.autoconfigure.amqp.RabbitAutoConfiguration
   data:
     mongodb:
       uri: mongodb://root:Yf%40123456@192.168.1.228:57017/com_yunfeiyun_iot
@@ -141,7 +144,15 @@ spring:
       connection-timeout: 15000
       publisher-returns: true
       enabled: true
-
+    agmpIot:
+      host: 192.168.1.230
+      port: 5672
+      username: admin
+      password: admin
+      virtual-host: /agmp-saas
+      connection-timeout: 15000
+      publisher-returns: true
+      enabled: true
 
   # redis 配置
   redis:
@@ -280,4 +291,4 @@ map:
 
 #只允许在228机器上访问,其他环境不要配置
 warn:
-  cbdDateDiff: 0
+  cbdDateDiff: 0

+ 0 - 264
src/main/resources/application-prod.yml

@@ -1,264 +0,0 @@
-# 项目相关配置
-application:
-  # 名称
-  name: IOTS
-  # 版本
-  version: 1.0.0
-  # 版权年份
-  copyrightYear: 2023
-  # 实例演示开关
-  demoEnabled: true
-  # 文件路径 示例( Windows配置D:/yunfei/farmwork/uploadPath,Linux配置 /home/yunfei/farmwork/uploadPath)
-  profile: /data/AGMP/iots
-  # 获取ip地址开关
-  addressEnabled: true
-  # 验证码类型 math 数组计算 char 字符验证
-  captchaType: math
-  # 顶级菜单的父Id
-  topMenuparentid: 0 #9a9d5a42-0803-5761-454c-177834ea4408
-  # 顶级部门的父Id
-  topDeptparentid: 8d4b26eb-1811-17a6-8dd1-127ea32c7cf8
-# 开发环境配置
-server:
-  # 服务器的HTTP端口
-  port: 8035
-  servlet:
-    # 应用的访问路径
-    context-path: /
-  tomcat:
-    # tomcat的URI编码
-    uri-encoding: UTF-8
-    # 连接数满后的排队数,默认为100
-    accept-count: 1000
-    threads:
-      # tomcat最大线程数,默认为200
-      max: 800
-      # Tomcat启动初始化的线程数,默认值10
-      min-spare: 100
-
-# 日志配置
-logging:
-  level:
-    com.yunfeiyun: debug
-    org.springframework: warn
-
-# 用户配置
-user:
-  password:
-    # 密码最大错误次数
-    maxRetryCount: 5
-    # 密码锁定时间(默认10分钟)
-    lockTime: 10
-
-# Spring配置
-spring:
-  # 资源信息
-  messages:
-    # 国际化资源文件路径
-    basename: i18n/messages
-  datasource:
-    type: com.alibaba.druid.pool.DruidDataSource
-    driverClassName: com.mysql.cj.jdbc.Driver
-    druid:
-      # 主库数据源
-      master:
-        url: jdbc:mysql://localhost:53306/com_yunfeiyun_agmp?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&useSSL=true&serverTimezone=GMT%2B8&allowMultiQueries=true
-        username: root
-        password: Yf@YqoGj#oG23
-      # 从库数据源
-      slave:
-        # 从数据源开关/默认关闭
-        enabled: false
-        url:
-        username:
-        password:
-      # 初始连接数
-      initialSize: 5
-      # 最小连接池数量
-      minIdle: 10
-      # 最大连接池数量
-      maxActive: 20
-      # 配置获取连接等待超时的时间
-      maxWait: 60000
-      # 配置连接超时时间
-      connectTimeout: 30000
-      # 配置网络超时时间
-      socketTimeout: 60000
-      # 配置间隔多久才进行一次检测,检测需要关闭的空闲连接,单位是毫秒
-      timeBetweenEvictionRunsMillis: 60000
-      # 配置一个连接在池中最小生存的时间,单位是毫秒
-      minEvictableIdleTimeMillis: 300000
-      # 配置一个连接在池中最大生存的时间,单位是毫秒
-      maxEvictableIdleTimeMillis: 900000
-      # 配置检测连接是否有效
-      validationQuery: SELECT 1 FROM DUAL
-      testWhileIdle: true
-      testOnBorrow: false
-      testOnReturn: false
-      webStatFilter:
-        enabled: true
-      statViewServlet:
-        enabled: true
-        # 设置白名单,不填则允许所有访问
-        allow:
-        url-pattern: /druid/*
-        # 控制台管理用户名和密码
-        login-username: ruoyi
-        login-password: 123456
-      filter:
-        stat:
-          enabled: true
-          # 慢SQL记录
-          log-slow-sql: true
-          slow-sql-millis: 1000
-          merge-sql: true
-        wall:
-          config:
-            multi-statement-allow: true
-  # 文件上传
-  servlet:
-    multipart:
-      # 单个文件大小
-      max-file-size:  10MB
-      # 设置总上传的文件大小
-      max-request-size:  20MB
-  # 服务模块
-  devtools:
-    restart:
-      # 热部署开关
-      enabled: true
-  rabbitmq:
-    host: localhost
-    port: 55763
-    username: yfkj_yanshi
-    password: Yf@TTki7F_katvL9
-    virtual-host: /agmpbs
-    connection-timeout: 15000
-    publisher-returns: true
-    enabled: true
-  # redis 配置
-  redis:
-    # 地址
-    host: localhost
-    # 端口,默认为6379
-    port: 56379
-    # 数据库索引
-    database: 0
-    # 密码
-    password: Yf@hsq9yVxp
-    # 连接超时时间
-    timeout: 10s
-    lettuce:
-      pool:
-        # 连接池中的最小空闲连接
-        min-idle: 0
-        # 连接池中的最大空闲连接
-        max-idle: 8
-        # 连接池的最大数据库连接数
-        max-active: 8
-        # #连接池最大阻塞等待时间(使用负值表示没有限制)
-        max-wait: -1ms
-
-# token配置
-token:
-  # 令牌自定义标识
-  header: Authorization
-  # 令牌密钥
-  secret: abcdefghijklmnopqrstuvwxyz
-  # 令牌有效期(默认30分钟)
-  expireTime: 30
-
-# MyBatis配置
-mybatis:
-  # 搜索指定包别名
-  typeAliasesPackage: com.yunfeiyun.**.domain
-  # 配置mapper的扫描,找到所有的mapper.xml映射文件
-  mapperLocations: classpath*:mapper/**/*Mapper.xml
-  # 加载全局的配置文件
-  configLocation: classpath:mybatis/mybatis-config.xml
-
-# PageHelper分页插件
-pagehelper:
-  helperDialect: mysql
-  supportMethodsArguments: true
-  params: count=countSql
-
-# Swagger配置
-swagger:
-  # 是否开启swagger
-  enabled: true
-  # 请求前缀
-  pathMapping: /dev-api
-
-# 防止XSS攻击
-xss:
-  # 过滤开关
-  enabled: true
-  # 排除链接(多个用逗号分隔)
-  excludes: /system/notice
-  # 匹配链接
-  urlPatterns: /system/*,/tool/*
-
-network:
-  ipaddrLogin: http://121.40.180.217
-  ipaddrSso: http://localhost
-  vueport: 7000
-  ssoport: 9002
-  icsport: 8023
-  wprport: 8027
-  fmsIp: localhost
-  fmsPort: 8021
-
-portal:
-  ssoEnabled: true
-  ssoLoginUrl: ${network.ipaddrLogin}:${network.vueport}/portal/login?redirect=
-  ssoTokenUrl: ${network.ipaddrSso}:${network.ssoport}/sso/auth/verifyToken
-  ssoSessionUrl: ${network.ipaddrSso}:${network.ssoport}/sso/auth/registerSession
-  ssoLogoutUrl: ${network.ipaddrSso}:${network.ssoport}/sso/auth/logout
-  appLogoutUrl: ${network.ipaddrSso}:${network.icsport}/sso/logoutToken
-  syncEnabled: true
-  orgEnabled: true
-
-policy:
-  # localSpace(本地空间), cloud 资源服务器
-  upload:
-    uploadType: localSpace
-    ossconfig:
-      ossType: 1
-      ossCloud:
-        aliYun:
-          accessType: 1
-          aliyunDomain: https://yunfei-agm.oss-cn-hangzhou.aliyuncs.com
-          aliyunPrefix: agmptest
-          aliyunEndPoint: oss-cn-hangzhou.aliyuncs.com
-          aliyunAccessKeyId: LTAI4G7tFh5Nk4KXZoSPk1D8
-          aliyunAccessKeySecret: RV4S2SfbLPoFNjlI4uIOoA0J1LQPQc
-          aliyunBucketName: yunfei-agm
-        qCloud:
-          qcloudDomain:
-          qcloudPrefix:
-          qcloudSecretId:
-          qcloudSecretKey:
-          qcloudBucketName: yunfei-agm
-  # table(表), cache(缓存)
-  water: table
-  queue:
-    productRouter:
-    consumeRouter:
-    ackEnabled: true
-iot:
-  customerId: 1
-  sassAble: false
-  ai:
-    rt:
-      callBackUrl: http://114.55.0.7:7000/iotsprod-api/ai/mc/camera/subscribe/image/callback
-
-weather:
-  api: http://open.nyzhwlw.com:10001/yf_weather
-  username: yunfeisaas
-  password: yf@yunfeisaas
-
-map:
-  gaode:
-    api: http://restapi.amap.com
-    key: 78ce288400f4fc6d9458989875c833c2

+ 14 - 0
src/main/resources/application-test.yml

@@ -52,6 +52,9 @@ user:
 
 # Spring配置
 spring:
+  autoconfigure:
+    exclude:
+      - org.springframework.boot.autoconfigure.amqp.RabbitAutoConfiguration
   data:
     mongodb:
       uri: mongodb://root:123456@127.0.0.1:27017/com_yunfeiyun_iot
@@ -131,11 +134,22 @@ spring:
       # 热部署开关
       enabled: true
   rabbitmq:
+    # 在物联网系统中 agmp 是 iotm 和 iots 之间通信使用
     agmp:
       host: localhost
       port: 5672
       username: user
       password: 123456
+      virtual-host: /agmp-saas-iot
+      connection-timeout: 15000
+      publisher-returns: true
+      enabled: true
+  # 物联网与门户之间使用
+    agmpIot:
+      host: localhost
+      port: 5672
+      username: user
+      password: 123456
       virtual-host: /agmp-saas
       connection-timeout: 15000
       publisher-returns: true

+ 55 - 0
src/main/resources/mapper/IotSfElementfactorMapper.xml

@@ -0,0 +1,55 @@
+<?xml version="1.0" encoding="UTF-8" ?>
+<!DOCTYPE mapper
+        PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
+        "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
+<mapper namespace="com.yunfeiyun.agmp.iots.mapper.IotSfElementfactorMapper">
+
+    <select id="selectIotSfElementfactorList" parameterType="IotSfElementfactorListReqVo"
+            resultType="com.yunfeiyun.agmp.iot.common.domain.IotSfElementfactor" >
+        SELECT sf.id, sf.sfBid, sf.devBid, sf.sfType, sf.sfCode, sf.sfName, sf.sfDisplayname, sf.sfParentBid, sf.sfSequence,
+            sf.tid, sf.sfCreatedDate, sf.sfCreator, sf.sfModifieddate, sf.sfModifier
+        FROM IotSfElementfactor AS sf
+        <where>
+            sf.tid = #{tid}
+            <if test="sfBid!= null and sfBid!= ''"> and sf.sfBid = #{sfBid}</if>
+            <if test="sfType!= null and sfType!= ''"> and sf.sfType = #{sfType}</if>
+            <if test="devBid != null and devBid != ''"> and sf.devBid = #{devBid}</if>
+            <if test="sfCode!= null and sfCode!= ''"> and sf.sfCode = #{sfCode}</if>
+            <if test="sfName!= null and sfName!= ''"> and sf.sfName like concat('%', #{sfName}, '%')</if>
+            <if test="sfDisplayname!= null and sfDisplayname!= ''"> and sf.sfDisplayname like concat('%', #{sfDisplayname}, '%')</if>
+            <if test="sfParentBid!= null and sfParentBid!= ''"> and sf.sfParentBid = #{sfParentBid}</if>
+            <if test="sfSequence!= null and sfSequence!= ''"> and sf.sfSequence = #{sfSequence}</if>
+            <if test="sfCreatedDate!= null and sfCreatedDate!= ''"> and sf.sfCreatedDate = #{sfCreatedDate}</if>
+            <if test="sfCreator!= null and sfCreator!= ''"> and sf.sfCreator = #{sfCreator}</if>
+            <if test="sfModifieddate!= null and sfModifieddate!= ''"> and sf.sfModifieddate = #{sfModifieddate}</if>
+            <if test="sfModifier!= null and sfModifier!= ''"> and sf.sfModifier = #{sfModifier}</if>
+            <if test="sfTypeList != null and sfTypeList.size() > 0">
+                and sf.sfType in
+                <foreach collection="sfTypeList" item="item" index="index" open="(" separator="," close=")">
+                    #{item}
+                </foreach>
+            </if>
+        </where>
+    </select>
+
+    <insert id="batchInsertIotSfElementfactor" parameterType="IotSfElementfactor">
+        INSERT INTO IotSfElementfactor (id, sfBid, devBid, sfType, sfCode, sfName, sfDisplayname, sfParentBid, sfSequence,
+            tid, sfCreatedDate, sfCreator, sfModifieddate, sfModifier)
+        VALUES
+        <foreach collection="list" item="item" index="index" separator=",">
+            (#{item.id}, #{item.sfBid}, #{item.devBid}, #{item.sfType}, #{item.sfCode}, #{item.sfName}, #{item.sfDisplayname}, #{item.sfParentBid}, #{item.sfSequence},
+            #{item.tid}, #{item.sfCreatedDate}, #{item.sfCreator}, #{item.sfModifieddate}, #{item.sfModifier})
+        </foreach>
+    </insert>
+
+    <delete id="deleteIotSfElementfactorBySfBid" parameterType="string">
+        DELETE FROM IotSfElementfactor WHERE sfBid = #{sfBid} OR sfParentBid = #{sfBid}
+    </delete>
+
+    <delete id="batchDeleteIotSfElementfactorBySfBidList" parameterType="string">
+        <foreach collection="list" item="item" index="index" separator=";">
+            DELETE FROM IotSfElementfactor WHERE sfBid = #{item} OR sfParentBid = #{item}
+        </foreach>
+    </delete>
+
+</mapper>

+ 145 - 0
src/main/resources/mapper/IotSfIrrigationRecordMapper.xml

@@ -0,0 +1,145 @@
+<?xml version="1.0" encoding="UTF-8" ?>
+<!DOCTYPE mapper
+        PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
+        "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
+<mapper namespace="com.yunfeiyun.agmp.iots.mapper.IotSfIrrigationRecordMapper">
+
+    <insert id="insertIrrigationRecord" parameterType="com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord">
+        INSERT INTO IotSfIrrigationRecord
+            <trim prefix="(" suffix=")" suffixOverrides=",">
+                <if test="rcdBid != null">rcdBid,</if>
+                <if test="devBid != null">devBid,</if>
+                <if test="rcdContent != null">rcdContent,</if>
+                <if test="rcdGroupbid != null">rcdGroupbid,</if>
+                <if test="rcdGroupName!= null">rcdGroupName,</if>
+                <if test="rcdStatus!= null">rcdStatus,</if>
+                <if test="rcdMode!= null">rcdMode,</if>
+                <if test="sfdataBid!= null">sfdataBid,</if>
+                <if test="rcdFlow!= null">rcdFlow,</if>
+                <if test="rcdCreator!= null">rcdCreator,</if>
+                <if test="rcdStartdate!= null">rcdStartdate,</if>
+                <if test="rcdEnddate!= null">rcdEnddate,</if>
+                <if test="rcdCreatorName!= null">rcdCreatorName,</if>
+                <if test="rcdCreateddate!= null">rcdCreateddate,</if>
+                <if test="tid!= null">tid,</if>
+                <if test="rcdTime!= null">rcdTime,</if>
+                <if test="rcdUpdateddate!= null">rcdUpdateddate,</if>
+            </trim>
+            <trim prefix="values (" suffix=")" suffixOverrides=",">
+                <if test="rcdBid!= null">#{rcdBid},</if>
+                <if test="devBid!= null">#{devBid},</if>
+                <if test="rcdContent!= null">#{rcdContent},</if>
+                <if test="rcdGroupbid!= null">#{rcdGroupbid},</if>
+                <if test="rcdGroupName!= null">#{rcdGroupName},</if>
+                <if test="rcdStatus!= null">#{rcdStatus},</if>
+                <if test="rcdMode!= null">#{rcdMode},</if>
+                <if test="sfdataBid!= null">#{sfdataBid},</if>
+                <if test="rcdFlow!= null">#{rcdFlow},</if>
+                <if test="rcdCreator!= null">#{rcdCreator},</if>
+                <if test="rcdStartdate!= null">#{rcdStartdate},</if>
+                <if test="rcdEnddate!= null">#{rcdEnddate},</if>
+                <if test="rcdCreatorName!= null">#{rcdCreatorName},</if>
+                <if test="rcdCreateddate!= null">#{rcdCreateddate},</if>
+                <if test="tid!= null">#{tid},</if>
+                <if test="rcdTime!= null">#{rcdTime},</if>
+                <if test="rcdFertilizer!= null">#{rcdFertilizer},</if>
+            </trim>
+    </insert>
+    <insert id="batchInsertIotSfIrrigationRecord">
+        <foreach collection="list" item="item" separator=";">
+            INSERT INTO IotSfIrrigationRecord
+            <trim prefix="(" suffix=")" suffixOverrides=",">
+                <if test="item.rcdBid!= null">rcdBid,</if>
+                <if test="item.devBid!= null">devBid,</if>
+                <if test="item.rcdContent!= null">rcdContent,</if>
+                <if test="item.rcdGroupbid!= null">rcdGroupbid,</if>
+                <if test="item.rcdGroupName!= null">rcdGroupName,</if>
+                <if test="item.rcdStatus!= null">rcdStatus,</if>
+                <if test="item.rcdMode!= null">rcdMode,</if>
+                <if test="item.sfdataBid!= null">sfdataBid,</if>
+                <if test="item.rcdFlow!= null">rcdFlow,</if>
+                <if test="item.rcdCreator!= null">rcdCreator,</if>
+                <if test="item.rcdStartdate!= null">rcdStartdate,</if>
+                <if test="item.rcdEnddate!= null">rcdEnddate,</if>
+                <if test="item.rcdCreatorName!= null">rcdCreatorName,</if>
+                <if test="item.rcdCreateddate!= null">rcdCreateddate,</if>
+                <if test="item.tid!= null">tid,</if>
+                <if test="item.rcdTime!= null">rcdTime,</if>
+                <if test="item.rcdFertilizer!= null">rcdFertilizer,</if>
+            </trim>
+            <trim prefix="values (" suffix=")" suffixOverrides=",">
+                <if test="item.rcdBid!= null">#{item.rcdBid},</if>
+                <if test="item.devBid!= null">#{item.devBid},</if>
+                <if test="item.rcdContent!= null">#{item.rcdContent},</if>
+                <if test="item.rcdGroupbid!= null">#{item.rcdGroupbid},</if>
+                <if test="item.rcdGroupName!= null">#{item.rcdGroupName},</if>
+                <if test="item.rcdStatus!= null">#{item.rcdStatus},</if>
+                <if test="item.rcdMode!= null">#{item.rcdMode},</if>
+                <if test="item.sfdataBid!= null">#{item.sfdataBid},</if>
+                <if test="item.rcdFlow!= null">#{item.rcdFlow},</if>
+                <if test="item.rcdCreator!= null">#{item.rcdCreator},</if>
+                <if test="item.rcdStartdate!= null">#{item.rcdStartdate},</if>
+                <if test="item.rcdEnddate!= null">#{item.rcdEnddate},</if>
+                <if test="item.rcdCreatorName!= null">#{item.rcdCreatorName},</if>
+                <if test="item.rcdCreateddate!= null">#{item.rcdCreateddate},</if>
+                <if test="item.tid!= null">#{item.tid},</if>
+                <if test="item.rcdTime!= null">#{item.rcdTime},</if>
+                <if test="item.rcdFertilizer!= null">#{item.rcdFertilizer},</if>
+            </trim>
+        </foreach>
+    </insert>
+
+    <update id="updateIrrigationRecord" parameterType="com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord">
+        UPDATE IotSfIrrigationRecord
+        <trim prefix="SET" suffixOverrides=",">
+            <if test="rcdContent!= null">rcdContent = #{rcdContent},</if>
+            <if test="rcdStatus!= null">rcdStatus = #{rcdStatus},</if>
+            <if test="rcdFlow!= null">rcdFlow = #{rcdFlow},</if>
+            <if test="rcdEnddate!= null">rcdEnddate = #{rcdEnddate},</if>
+            <if test="rcdTime!= null">rcdTime = #{rcdTime},</if>
+            <if test="rcdFertilizer!= null">rcdFertilizer = #{rcdFertilizer},</if>
+        </trim>
+        where rcdBid = #{rcdBid} and tid = #{tid}
+    </update>
+
+    <update id="batchUpdateIrrigationRecord" parameterType="IotSfIrrigationRecord">
+        <foreach collection="list" item="item" separator=";">
+            UPDATE IotSfIrrigationRecord
+            <trim prefix="SET" suffixOverrides=",">
+                <if test="item.rcdContent!= null">rcdContent = #{item.rcdContent},</if>
+                <if test="item.rcdStatus!= null">rcdStatus = #{item.rcdStatus},</if>
+                <if test="item.rcdFlow!= null">rcdFlow = #{item.rcdFlow},</if>
+                <if test="item.rcdEnddate!= null">rcdEnddate = #{item.rcdEnddate},</if>
+                <if test="item.rcdTime!= null">rcdTime = #{item.rcdTime},</if>
+                <if test="item.rcdFertilizer!= null">rcdFertilizer = #{item.rcdFertilizer},</if>
+            </trim>
+            where rcdBid = #{item.rcdBid} and tid = #{item.tid}
+        </foreach>
+    </update>
+
+    <select id="selectIrrigationRecordByBid" parameterType="com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord"
+            resultType="com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord">
+        SELECT rcdBid, devBid, rcdContent, rcdGroupbid, rcdGroupName, rcdStatus, rcdMode, sfdataBid, rcdFlow, rcdCreator,
+            rcdStartdate, rcdEnddate, rcdCreatorName, rcdCreateddate, tid, rcdTime, rcdFertilizer
+        FROM IotSfIrrigationRecord
+        where rcdBid = #{rcdBid}
+    </select>
+
+    <select id="selectIrrigationRecordList" parameterType="IotSfIrrigationRecordListReqVo"
+            resultType="com.yunfeiyun.agmp.iot.common.domain.IotSfIrrigationRecord">
+        SELECT rcdBid, devBid, rcdContent, rcdGroupbid, rcdGroupName, rcdStatus, rcdMode, sfdataBid, rcdFlow, rcdCreator,
+            rcdStartdate, rcdEnddate, rcdCreatorName, rcdCreateddate, tid, rcdTime, rcdFertilizer
+        FROM IotSfIrrigationRecord
+        <where>
+            tid = #{tid}
+            <if test="devBid!= null and devBid!= ''"> and devBid = #{devBid}</if>
+            <if test="rcdGroupbid!= null and rcdGroupbid!= ''"> and rcdGroupbid = #{rcdGroupbid}</if>
+            <if test="rcdGroupName!= null and rcdGroupName!= ''"> and rcdGroupName = #{rcdGroupName}</if>
+            <if test="rcdStatus!= null and rcdStatus!= ''"> and rcdStatus = #{rcdStatus}</if>
+            <if test="rcdMode!= null and rcdMode!= ''"> and rcdMode = #{rcdMode}</if>
+            <if test="startTime!= null and startTime!= ''"> and rcdStartdate <![CDATA[ >= ]]> #{startTime}</if>
+            <if test="endTime!= null and endTime!= ''"> and rcdStartdate <![CDATA[ <= ]]> #{endTime}</if>
+        </where>
+        order by rcdCreateddate desc
+    </select>
+</mapper>

+ 38 - 0
src/main/resources/mapper/IotWarnBusinessMapper.xml

@@ -22,6 +22,7 @@
             <if test="wlData != null">wlData,</if>
             <if test="tid != null and tid != ''">tid,</if>
             <if test="wcBid != null and wcBid != ''">wcBid,</if>
+            <if test="wlSendmsgstatus != null and wlSendmsgstatus != ''">wlSendmsgstatus,</if>
         </trim>
         <trim prefix="values (" suffix=")" suffixOverrides=",">
             <if test="wlBid != null">#{wlBid},</if>
@@ -39,6 +40,7 @@
             <if test="wlData != null">#{wlData},</if>
             <if test="tid != null and tid != ''">#{tid},</if>
             <if test="wcBid != null and wcBid != ''">#{wcBid},</if>
+            <if test="wlSendmsgstatus != null and wlSendmsgstatus != ''">#{wlSendmsgstatus},</if>
         </trim>
     </insert>
     <insert id="insertIncrementReCount" parameterType="IotWarncount" useGeneratedKeys="true" keyProperty="id">
@@ -167,6 +169,9 @@
         where devBid = #{devBid} and wlType='1' and status='0'
 
     </update>
+    <update id="updateWarnLogSendStatus">
+        update IotWarnlog set wlSendmsgstatus=#{status},wlSendmsgtime=#{time}  where wlBid=#{wlBid}
+    </update>
 
     <select id="selectIotWarnconfigCbdDevList" resultType="com.yunfeiyun.agmp.iots.warn.model.IotWarnconfigDevVo">
         SELECT wi.wiBid, d.devBid, d.devCode, d.devtypeBid, d.devCbdrecogtype, wc.*
@@ -216,5 +221,38 @@
             </if>
         </where>
     </select>
+    <select id="selectYbqIndicatorAllList" resultType="com.yunfeiyun.agmp.iot.common.domain.IotWarnindicator">
+
 
+        SELECT wi.*
+        FROM IotWarnindicator AS wi
+        LEFT JOIN IotWarnconfig AS wc ON wc.wcBid = wi.wcBid
+        WHERE wi.wiCode IN ('computeValue') AND wc.wcStatus = "0"
+
+    </select>
+
+    <select id="selectIotWarnconfigYbqDevList"
+            resultType="com.yunfeiyun.agmp.iots.warn.model.IotWarnconfigDevVo">
+                  SELECT wi.wiBid, d.devBid, d.devCode, d.devtypeBid, d.devCbdrecogtype, wc.*
+        FROM IotWarnindicator AS wi
+            LEFT JOIN IotWarnconfig AS wc ON wc.wcBid = wi.wcBid
+            LEFT JOIN IotWarnobject AS wo ON wo.wcBid = wc.wcBid
+            LEFT JOIN IotDevice AS d ON d.devBid = wo.devBid
+        WHERE wi.wiCode = "computeValue" AND wc.wcStatus = "0" AND d.devDelstatus = "0"
+    </select>
+    <select id="selectDeviceById" resultType="com.yunfeiyun.agmp.iot.common.domain.IotDevice">
+        select * from IotDevice where devBid=#{devBid}
+    </select>
+    <select id="selectWarnPolicy" resultType="com.yunfeiyun.agmp.iot.common.domain.IotWarnpolicy">
+        select *
+        from IotWarnpolicy
+        where wcBid = #{configId}
+    </select>
+    <select id="selectWarnReceiverByConfigId"
+            resultType="com.yunfeiyun.agmp.iot.common.domain.IotWarnreceiver">
+        select * from IotWarnreceiver where wcBid = #{configId}
+    </select>
+    <select id="getLastedUnSendWarnLog" resultType="com.yunfeiyun.agmp.iot.common.domain.IotWarnlog">
+        select * from IotWarnlog where wcBid = #{wcBid} and wlSendmsgstatus = "1" order by wlCreateddate desc limit 1
+    </select>
 </mapper>