소스 검색

mqtt收数据转发及指令发布

liwei19941102 2 년 전
부모
커밋
bbf4472ca0

+ 12 - 9
ruoyi-common/src/main/java/com/ruoyi/common/utils/mqtt/MQTTConnect.java

@@ -28,16 +28,19 @@ public class MQTTConnect {
      * @param mqttCallback 回调函数
      * @param mqttCallback 回调函数
      **/
      **/
     public void createMqttClient(String Host,String clientId,String userName, String passWord, MqttCallback mqttCallback) throws MqttException {
     public void createMqttClient(String Host,String clientId,String userName, String passWord, MqttCallback mqttCallback) throws MqttException {
-        MqttConnectOptions options = mqttConnectOptions(Host,clientId,userName, passWord);
-        if (mqttCallback == null) {
-            mqttClient.setCallback(new Callback());
-        } else {
-            mqttClient.setCallback(mqttCallback);
+        mqttClient = (MqttClient) CacheManager.getCacheDataByKey("mqtt"+clientId);
+        if(mqttClient == null){
+            MqttConnectOptions options = mqttConnectOptions(Host,clientId,userName, passWord);
+            if (mqttCallback == null) {
+                mqttClient.setCallback(new Callback());
+            } else {
+                mqttClient.setCallback(mqttCallback);
+            }
+            mqttClient.connect(options);
+            String key = "mqtt"+clientId;
+            CacheManagerEntity cacheManagerEntity = new CacheManagerEntity(mqttClient);
+            CacheManager.putCache(key,cacheManagerEntity);
         }
         }
-        mqttClient.connect(options);
-        String key = "mqtt"+clientId;
-        CacheManagerEntity cacheManagerEntity = new CacheManagerEntity(mqttClient);
-        CacheManager.putCache(key,cacheManagerEntity);
     }
     }
 
 
     /**
     /**

+ 59 - 0
ruoyi-system/src/main/java/com/ruoyi/data/controller/MqttController.java

@@ -0,0 +1,59 @@
+package com.ruoyi.data.controller;
+
+import cn.dev33.satoken.annotation.SaCheckPermission;
+import com.ruoyi.common.annotation.Log;
+import com.ruoyi.common.enums.BusinessType;
+import com.ruoyi.common.utils.mqtt.MQTTConnect;
+import com.ruoyi.data.domain.OrderBean;
+import com.ruoyi.data.domain.TblMqtt;
+import com.ruoyi.data.domain.bo.TblMqttBo;
+import com.ruoyi.data.domain.vo.TblMqttVo;
+import com.ruoyi.data.service.ITblMqttService;
+import com.ruoyi.data.service.MqttService;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
+import org.eclipse.paho.client.mqttv3.MqttCallback;
+import org.eclipse.paho.client.mqttv3.MqttClient;
+import org.eclipse.paho.client.mqttv3.MqttMessage;
+import org.springframework.integration.mqtt.support.MqttUtils;
+import org.springframework.validation.annotation.Validated;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import java.util.List;
+
+@Validated
+@RequiredArgsConstructor
+@RestController
+@RequestMapping("/data/mqtt")
+@Slf4j
+public class MqttController {
+
+
+//    private final ITblMqttService mqttService;
+
+    private final MqttService mqttService;
+
+    @GetMapping("/createMqtt")
+    public void createMqtt(TblMqttBo bo){
+
+    }
+
+    @GetMapping("/createMqttMain")
+    public void createMqttMain(TblMqttBo bo){
+        mqttService.createMqttMain(bo);
+    }
+
+    @GetMapping("/pubOrder")
+    public void pubOrder(OrderBean bean){
+        mqttService.pubOrder(bean);
+    }
+
+
+
+}
+
+

+ 108 - 0
ruoyi-system/src/main/java/com/ruoyi/data/controller/TblRecordController.java

@@ -0,0 +1,108 @@
+package com.ruoyi.data.controller;
+
+import java.util.List;
+import java.util.Arrays;
+import java.util.concurrent.TimeUnit;
+
+import lombok.RequiredArgsConstructor;
+import javax.servlet.http.HttpServletResponse;
+import javax.validation.constraints.*;
+import cn.dev33.satoken.annotation.SaCheckPermission;
+import org.springframework.web.bind.annotation.*;
+import org.springframework.validation.annotation.Validated;
+import com.ruoyi.common.annotation.RepeatSubmit;
+import com.ruoyi.common.annotation.Log;
+import com.ruoyi.common.core.controller.BaseController;
+import com.ruoyi.common.core.domain.PageQuery;
+import com.ruoyi.common.core.domain.R;
+import com.ruoyi.common.core.validate.AddGroup;
+import com.ruoyi.common.core.validate.EditGroup;
+import com.ruoyi.common.core.validate.QueryGroup;
+import com.ruoyi.common.enums.BusinessType;
+import com.ruoyi.common.utils.poi.ExcelUtil;
+import com.ruoyi.data.domain.vo.TblRecordVo;
+import com.ruoyi.data.domain.bo.TblRecordBo;
+import com.ruoyi.data.service.ITblRecordService;
+import com.ruoyi.common.core.page.TableDataInfo;
+
+/**
+ * 数据记录
+ *
+ * @author ruoyi
+ * @date 2023-12-14
+ */
+@Validated
+@RequiredArgsConstructor
+@RestController
+@RequestMapping("/data/record")
+public class TblRecordController extends BaseController {
+
+    private final ITblRecordService iTblRecordService;
+
+    /**
+     * 查询数据记录列表
+     */
+    @SaCheckPermission("data:record:list")
+    @GetMapping("/list")
+    public TableDataInfo<TblRecordVo> list(TblRecordBo bo, PageQuery pageQuery) {
+        return iTblRecordService.queryPageList(bo, pageQuery);
+    }
+
+    /**
+     * 导出数据记录列表
+     */
+    @SaCheckPermission("data:record:export")
+    @Log(title = "数据记录", businessType = BusinessType.EXPORT)
+    @PostMapping("/export")
+    public void export(TblRecordBo bo, HttpServletResponse response) {
+        List<TblRecordVo> list = iTblRecordService.queryList(bo);
+        ExcelUtil.exportExcel(list, "数据记录", TblRecordVo.class, response);
+    }
+
+    /**
+     * 获取数据记录详细信息
+     *
+     * @param id 主键
+     */
+    @SaCheckPermission("data:record:query")
+    @GetMapping("/{id}")
+    public R<TblRecordVo> getInfo(@NotNull(message = "主键不能为空")
+                                     @PathVariable Long id) {
+        return R.ok(iTblRecordService.queryById(id));
+    }
+
+    /**
+     * 新增数据记录
+     */
+    @SaCheckPermission("data:record:add")
+    @Log(title = "数据记录", businessType = BusinessType.INSERT)
+    @RepeatSubmit()
+    @PostMapping()
+    public R<Void> add(@Validated(AddGroup.class) @RequestBody TblRecordBo bo) {
+        return toAjax(iTblRecordService.insertByBo(bo));
+    }
+
+    /**
+     * 修改数据记录
+     */
+    @SaCheckPermission("data:record:edit")
+    @Log(title = "数据记录", businessType = BusinessType.UPDATE)
+    @RepeatSubmit()
+    @PutMapping()
+    public R<Void> edit(@Validated(EditGroup.class) @RequestBody TblRecordBo bo) {
+        return toAjax(iTblRecordService.updateByBo(bo));
+    }
+
+    /**
+     * 删除数据记录
+     *
+     * @param ids 主键串
+     */
+    @SaCheckPermission("data:record:remove")
+    @Log(title = "数据记录", businessType = BusinessType.DELETE)
+    @DeleteMapping("/{ids}")
+    public R<Void> remove(@NotEmpty(message = "主键不能为空")
+                          @PathVariable Long[] ids) {
+        return toAjax(iTblRecordService.deleteWithValidByIds(Arrays.asList(ids), true));
+    }
+}

+ 15 - 0
ruoyi-system/src/main/java/com/ruoyi/data/domain/OrderBean.java

@@ -0,0 +1,15 @@
+package com.ruoyi.data.domain;
+
+import lombok.Data;
+
+@Data
+public class OrderBean {
+
+    private String value;
+
+    private Integer add;
+
+    private Integer addrOffset;
+
+    private Long deviceId;
+}

+ 42 - 0
ruoyi-system/src/main/java/com/ruoyi/data/domain/TblRecord.java

@@ -0,0 +1,42 @@
+package com.ruoyi.data.domain;
+
+import com.baomidou.mybatisplus.annotation.*;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import java.io.Serializable;
+import java.util.Date;
+import java.math.BigDecimal;
+
+import com.ruoyi.common.core.domain.BaseEntity;
+
+/**
+ * 数据记录对象 tbl_record
+ *
+ * @author ruoyi
+ * @date 2023-12-14
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@TableName("tbl_record")
+public class TblRecord extends BaseEntity {
+
+    private static final long serialVersionUID=1L;
+
+    /**
+     *
+     */
+    @TableId(value = "id")
+    private Long id;
+    /**
+     *
+     */
+    private String json;
+    /**
+     *
+     */
+    /**
+     *
+     */
+    private Long equipmentId;
+
+}

+ 44 - 0
ruoyi-system/src/main/java/com/ruoyi/data/domain/bo/TblRecordBo.java

@@ -0,0 +1,44 @@
+package com.ruoyi.data.domain.bo;
+
+import com.ruoyi.common.core.validate.AddGroup;
+import com.ruoyi.common.core.validate.EditGroup;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import javax.validation.constraints.*;
+
+import java.util.Date;
+
+import com.ruoyi.common.core.domain.BaseEntity;
+
+/**
+ * 数据记录业务对象 tbl_record
+ *
+ * @author ruoyi
+ * @date 2023-12-14
+ */
+
+@Data
+@EqualsAndHashCode(callSuper = true)
+public class TblRecordBo extends BaseEntity {
+
+    /**
+     *
+     */
+    @NotNull(message = "不能为空", groups = { EditGroup.class })
+    private Long id;
+
+    /**
+     *
+     */
+    @NotBlank(message = "不能为空", groups = { AddGroup.class, EditGroup.class })
+    private String json;
+
+
+    /**
+     *
+     */
+    @NotNull(message = "不能为空", groups = { AddGroup.class, EditGroup.class })
+    private Long equipmentId;
+
+
+}

+ 44 - 0
ruoyi-system/src/main/java/com/ruoyi/data/domain/vo/TblRecordVo.java

@@ -0,0 +1,44 @@
+package com.ruoyi.data.domain.vo;
+
+import com.alibaba.excel.annotation.ExcelIgnoreUnannotated;
+import com.alibaba.excel.annotation.ExcelProperty;
+import com.ruoyi.common.annotation.ExcelDictFormat;
+import com.ruoyi.common.convert.ExcelDictConvert;
+import lombok.Data;
+import java.util.Date;
+
+import java.io.Serializable;
+
+/**
+ * 数据记录视图对象 tbl_record
+ *
+ * @author ruoyi
+ * @date 2023-12-14
+ */
+@Data
+@ExcelIgnoreUnannotated
+public class TblRecordVo implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    /**
+     *
+     */
+    @ExcelProperty(value = "")
+    private Long id;
+
+    /**
+     *
+     */
+    @ExcelProperty(value = "")
+    private String json;
+
+
+    /**
+     *
+     */
+    @ExcelProperty(value = "")
+    private Long equipmentId;
+
+
+}

+ 4 - 0
ruoyi-system/src/main/java/com/ruoyi/data/mapper/TblEquipmentMqttMapper.java

@@ -1,9 +1,12 @@
 package com.ruoyi.data.mapper;
 package com.ruoyi.data.mapper;
 
 
+import com.ruoyi.data.domain.MqttObj;
 import com.ruoyi.data.domain.TblEquipmentMqtt;
 import com.ruoyi.data.domain.TblEquipmentMqtt;
 import com.ruoyi.data.domain.vo.TblEquipmentMqttVo;
 import com.ruoyi.data.domain.vo.TblEquipmentMqttVo;
 import com.ruoyi.common.core.mapper.BaseMapperPlus;
 import com.ruoyi.common.core.mapper.BaseMapperPlus;
 
 
+import java.util.List;
+
 /**
 /**
  * 【请填写功能名称】Mapper接口
  * 【请填写功能名称】Mapper接口
  *
  *
@@ -12,4 +15,5 @@ import com.ruoyi.common.core.mapper.BaseMapperPlus;
  */
  */
 public interface TblEquipmentMqttMapper extends BaseMapperPlus<TblEquipmentMqttMapper, TblEquipmentMqtt, TblEquipmentMqttVo> {
 public interface TblEquipmentMqttMapper extends BaseMapperPlus<TblEquipmentMqttMapper, TblEquipmentMqtt, TblEquipmentMqttVo> {
 
 
+    List<MqttObj> selectMqttListByDeviceId(MqttObj mqttObj);
 }
 }

+ 15 - 0
ruoyi-system/src/main/java/com/ruoyi/data/mapper/TblRecordMapper.java

@@ -0,0 +1,15 @@
+package com.ruoyi.data.mapper;
+
+import com.ruoyi.data.domain.TblRecord;
+import com.ruoyi.data.domain.vo.TblRecordVo;
+import com.ruoyi.common.core.mapper.BaseMapperPlus;
+
+/**
+ * 数据记录Mapper接口
+ *
+ * @author ruoyi
+ * @date 2023-12-14
+ */
+public interface TblRecordMapper extends BaseMapperPlus<TblRecordMapper, TblRecord, TblRecordVo> {
+
+}

+ 49 - 0
ruoyi-system/src/main/java/com/ruoyi/data/service/ITblRecordService.java

@@ -0,0 +1,49 @@
+package com.ruoyi.data.service;
+
+import com.ruoyi.data.domain.TblRecord;
+import com.ruoyi.data.domain.vo.TblRecordVo;
+import com.ruoyi.data.domain.bo.TblRecordBo;
+import com.ruoyi.common.core.page.TableDataInfo;
+import com.ruoyi.common.core.domain.PageQuery;
+
+import java.util.Collection;
+import java.util.List;
+
+/**
+ * 数据记录Service接口
+ *
+ * @author ruoyi
+ * @date 2023-12-14
+ */
+public interface ITblRecordService {
+
+    /**
+     * 查询数据记录
+     */
+    TblRecordVo queryById(Long id);
+
+    /**
+     * 查询数据记录列表
+     */
+    TableDataInfo<TblRecordVo> queryPageList(TblRecordBo bo, PageQuery pageQuery);
+
+    /**
+     * 查询数据记录列表
+     */
+    List<TblRecordVo> queryList(TblRecordBo bo);
+
+    /**
+     * 新增数据记录
+     */
+    Boolean insertByBo(TblRecordBo bo);
+
+    /**
+     * 修改数据记录
+     */
+    Boolean updateByBo(TblRecordBo bo);
+
+    /**
+     * 校验并批量删除数据记录信息
+     */
+    Boolean deleteWithValidByIds(Collection<Long> ids, Boolean isValid);
+}

+ 15 - 0
ruoyi-system/src/main/java/com/ruoyi/data/service/MqttService.java

@@ -0,0 +1,15 @@
+package com.ruoyi.data.service;
+
+import com.ruoyi.data.domain.OrderBean;
+import com.ruoyi.data.domain.bo.TblMqttBo;
+
+public interface MqttService {
+
+    void pubMqttData(String mqttStr);
+
+    void createMqttMain(TblMqttBo bo);
+
+    void createMqtt(TblMqttBo bo);
+
+    void pubOrder(OrderBean orderBean);
+}

+ 172 - 0
ruoyi-system/src/main/java/com/ruoyi/data/service/impl/MqttServiceImpl.java

@@ -0,0 +1,172 @@
+package com.ruoyi.data.service.impl;
+
+import cn.hutool.json.JSONObject;
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
+import com.ruoyi.common.utils.StringUtils;
+import com.ruoyi.common.utils.mqtt.MQTTConnect;
+import com.ruoyi.data.domain.MqttObj;
+import com.ruoyi.data.domain.OrderBean;
+import com.ruoyi.data.domain.TblMqtt;
+import com.ruoyi.data.domain.TblRecord;
+import com.ruoyi.data.domain.bo.TblMqttBo;
+import com.ruoyi.data.domain.vo.TblMqttVo;
+import com.ruoyi.data.mapper.TblEquipmentMqttMapper;
+import com.ruoyi.data.mapper.TblMqttMapper;
+import com.ruoyi.data.mapper.TblRecordMapper;
+import com.ruoyi.data.service.MqttService;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
+import org.eclipse.paho.client.mqttv3.MqttCallback;
+import org.eclipse.paho.client.mqttv3.MqttMessage;
+import org.springframework.stereotype.Service;
+
+import java.text.SimpleDateFormat;
+import java.util.Base64;
+import java.util.Date;
+import java.util.List;
+import java.util.Map;
+
+@RequiredArgsConstructor
+@Service
+@Slf4j
+public class MqttServiceImpl implements MqttService {
+
+    private final TblMqttMapper tblMqttMapper;
+
+    private final TblEquipmentMqttMapper tblEquipmentMqttMapper;
+
+    private final TblRecordMapper tblRecordMapper;
+
+
+    @Override
+    public void pubMqttData(String mqttStr) {
+        JSONObject jsonObject = new JSONObject(mqttStr);
+        Long deviceId = Long.valueOf((String) jsonObject.get("deviceId"));
+        SimpleDateFormat formatter= new SimpleDateFormat("yyyy-MM-dd 'at' HH:mm:ss");
+        Date date = new Date(System.currentTimeMillis());
+        jsonObject.put("created_time",formatter.format(date));
+        MqttObj mqttObj = new MqttObj();
+        mqttObj.setEquipmentId(deviceId);
+        List<MqttObj> mqttObjList = tblEquipmentMqttMapper.selectMqttListByDeviceId(mqttObj);
+        for(MqttObj obj:mqttObjList){
+                JSONObject topicObj = obj.getTopicQos("tcp");
+                MQTTConnect mqttConnect = new MQTTConnect();
+                try{
+                    mqttConnect.createMqttClient(obj.getServerAddress(),obj.getUuid(),obj.getAccount(),obj.getPassword(),new Callback());
+                    if(topicObj != null){
+                        String topic = topicObj.get("name").toString().replace("#","");
+                        mqttConnect.pub(topic,mqttStr,Integer.valueOf((String) topicObj.get("qos")));
+                    }else{
+                        String topic = "sensor/modbustcp/"+deviceId;
+                        mqttConnect.pub(topic,mqttStr,0);
+                    }
+                    TblRecord tblRecord = new TblRecord();
+                    tblRecord.setEquipmentId(deviceId);
+                    tblRecord.setJson(mqttStr);
+                    tblRecord.setCreateBy("admin");
+                    tblRecord.setUpdateBy("admin");
+                    tblRecordMapper.insert(tblRecord);
+                }catch (Exception e){
+                    e.printStackTrace();
+                }
+            }
+    }
+
+    @Override
+    public void createMqttMain(TblMqttBo bo){
+        MQTTConnect mqttConnect = new MQTTConnect();
+        try {
+            mqttConnect.createMqttClient("ws://52.130.249.112:8083/mqtt","adminTest","ship","ship@2021.11.24",new Callback());
+            mqttConnect.sub("device/#");
+        }catch (Exception e){
+            e.printStackTrace();
+        }
+
+    }
+
+    @Override
+    public void createMqtt(TblMqttBo bo){
+        LambdaQueryWrapper<TblMqtt> lqw = buildQueryWrapper(bo);
+        List<TblMqttVo> mqttVoList = tblMqttMapper.selectVoList(lqw);
+        for(TblMqttVo mqttVo:mqttVoList){
+            MQTTConnect mqttConnect = new MQTTConnect();
+            try {
+                mqttConnect.createMqttClient(mqttVo.getServerAddress(),mqttVo.getUuid(),mqttVo.getAccount(),mqttVo.getPassword(),new Callback());
+            }catch (Exception e){
+                e.printStackTrace();
+            }
+        }
+    }
+
+    @Override
+    public void pubOrder(OrderBean orderBean){
+        JSONObject jsonObject = new JSONObject();
+        jsonObject.put("deviceId",orderBean.getDeviceId());
+        jsonObject.put("add",orderBean.getAdd());
+        jsonObject.put("value",orderBean.getValue());
+        jsonObject.put("addrOffset",orderBean.getAddrOffset());
+        jsonObject.put("len",orderBean.getValue().split(",").length);
+        MQTTConnect mqttConnect = new MQTTConnect();
+        try {
+            mqttConnect.createMqttClient("ws://52.130.249.112:8083/mqtt","adminTest","ship","ship@2021.11.24",new Callback());
+            mqttConnect.pub("control",jsonObject.toString(),0);
+        }catch (Exception e){
+            e.printStackTrace();
+        }
+    }
+
+    private LambdaQueryWrapper<TblMqtt> buildQueryWrapper(TblMqttBo bo) {
+        Map<String, Object> params = bo.getParams();
+        LambdaQueryWrapper<TblMqtt> lqw = Wrappers.lambdaQuery();
+        lqw.like(StringUtils.isNotBlank(bo.getProtocolName()), TblMqtt::getProtocolName, bo.getProtocolName());
+        lqw.eq(StringUtils.isNotBlank(bo.getProtocolDesc()), TblMqtt::getProtocolDesc, bo.getProtocolDesc());
+        lqw.eq(StringUtils.isNotBlank(bo.getProtocolType()), TblMqtt::getProtocolType, bo.getProtocolType());
+        lqw.eq(StringUtils.isNotBlank(bo.getServerAddress()), TblMqtt::getServerAddress, bo.getServerAddress());
+        lqw.eq(StringUtils.isNotBlank(bo.getServerTopic()), TblMqtt::getServerTopic, bo.getServerTopic());
+        lqw.eq(StringUtils.isNotBlank(bo.getAccount()), TblMqtt::getAccount, bo.getAccount());
+        lqw.eq(StringUtils.isNotBlank(bo.getPassword()), TblMqtt::getPassword, bo.getPassword());
+        lqw.eq(StringUtils.isNotBlank(bo.getUuid()), TblMqtt::getUuid, bo.getUuid());
+//        lqw.eq(bo.getStatus() != null, TblMqtt::getStatus, bo.getStatus());
+        return lqw;
+    }
+
+    class Callback implements MqttCallback {
+
+        /**
+         * MQTT 断开连接会执行此方法
+         */
+        @Override
+        public void connectionLost(Throwable throwable) {
+            log.info("断开了MQTT连接111 :{}", throwable.getMessage());
+            log.error(throwable.getMessage(), throwable);
+        }
+
+        /**
+         * publish发布成功后会执行到这里
+         */
+        @Override
+        public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) {
+            log.info("发布消息成功");
+        }
+
+        /**
+         * subscribe订阅后得到的消息会执行到这里
+         */
+        @Override
+        public void messageArrived(String topic, MqttMessage message) throws Exception {
+            //  TODO    此处可以将订阅得到的消息进行业务处理、数据存储
+            String payload = String.valueOf(message.getPayload());
+           // String msg = Byte.toString(message.getPayload());
+            byte[] bytes = message.getPayload();
+            String encoded = Base64.getEncoder().encodeToString(bytes);
+            byte[] decoded = Base64.getDecoder().decode(encoded);
+            String msg =new String(decoded);
+            System.out.println(msg);
+            log.info("收到来自 " + topic + " 的消息:{}", new String(message.getPayload()));
+            MqttServiceImpl.this.pubMqttData(msg);
+        }
+    }
+
+}

+ 110 - 0
ruoyi-system/src/main/java/com/ruoyi/data/service/impl/TblRecordServiceImpl.java

@@ -0,0 +1,110 @@
+package com.ruoyi.data.service.impl;
+
+import cn.hutool.core.bean.BeanUtil;
+import com.ruoyi.common.utils.StringUtils;
+import com.ruoyi.common.core.page.TableDataInfo;
+import com.ruoyi.common.core.domain.PageQuery;
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
+import lombok.RequiredArgsConstructor;
+import org.springframework.stereotype.Service;
+import com.ruoyi.data.domain.bo.TblRecordBo;
+import com.ruoyi.data.domain.vo.TblRecordVo;
+import com.ruoyi.data.domain.TblRecord;
+import com.ruoyi.data.mapper.TblRecordMapper;
+import com.ruoyi.data.service.ITblRecordService;
+
+import java.util.List;
+import java.util.Map;
+import java.util.Collection;
+
+/**
+ * 数据记录Service业务层处理
+ *
+ * @author ruoyi
+ * @date 2023-12-14
+ */
+@RequiredArgsConstructor
+@Service
+public class TblRecordServiceImpl implements ITblRecordService {
+
+    private final TblRecordMapper baseMapper;
+
+    /**
+     * 查询数据记录
+     */
+    @Override
+    public TblRecordVo queryById(Long id){
+        return baseMapper.selectVoById(id);
+    }
+
+    /**
+     * 查询数据记录列表
+     */
+    @Override
+    public TableDataInfo<TblRecordVo> queryPageList(TblRecordBo bo, PageQuery pageQuery) {
+        LambdaQueryWrapper<TblRecord> lqw = buildQueryWrapper(bo);
+        Page<TblRecordVo> result = baseMapper.selectVoPage(pageQuery.build(), lqw);
+        return TableDataInfo.build(result);
+    }
+
+    /**
+     * 查询数据记录列表
+     */
+    @Override
+    public List<TblRecordVo> queryList(TblRecordBo bo) {
+        LambdaQueryWrapper<TblRecord> lqw = buildQueryWrapper(bo);
+        return baseMapper.selectVoList(lqw);
+    }
+
+    private LambdaQueryWrapper<TblRecord> buildQueryWrapper(TblRecordBo bo) {
+        Map<String, Object> params = bo.getParams();
+        LambdaQueryWrapper<TblRecord> lqw = Wrappers.lambdaQuery();
+        lqw.eq(StringUtils.isNotBlank(bo.getJson()), TblRecord::getJson, bo.getJson());
+        lqw.eq(bo.getEquipmentId() != null, TblRecord::getEquipmentId, bo.getEquipmentId());
+        return lqw;
+    }
+
+    /**
+     * 新增数据记录
+     */
+    @Override
+    public Boolean insertByBo(TblRecordBo bo) {
+        TblRecord add = BeanUtil.toBean(bo, TblRecord.class);
+        validEntityBeforeSave(add);
+        boolean flag = baseMapper.insert(add) > 0;
+        if (flag) {
+            bo.setId(add.getId());
+        }
+        return flag;
+    }
+
+    /**
+     * 修改数据记录
+     */
+    @Override
+    public Boolean updateByBo(TblRecordBo bo) {
+        TblRecord update = BeanUtil.toBean(bo, TblRecord.class);
+        validEntityBeforeSave(update);
+        return baseMapper.updateById(update) > 0;
+    }
+
+    /**
+     * 保存前的数据校验
+     */
+    private void validEntityBeforeSave(TblRecord entity){
+        //TODO 做一些数据校验,如唯一约束
+    }
+
+    /**
+     * 批量删除数据记录
+     */
+    @Override
+    public Boolean deleteWithValidByIds(Collection<Long> ids, Boolean isValid) {
+        if(isValid){
+            //TODO 做一些业务上的校验,判断是否需要校验
+        }
+        return baseMapper.deleteBatchIds(ids) > 0;
+    }
+}

+ 27 - 1
ruoyi-system/src/main/resources/mapper/data/TblEquipmentMqttMapper.xml

@@ -4,11 +4,37 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
 "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
 "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
 <mapper namespace="com.ruoyi.data.mapper.TblEquipmentMqttMapper">
 <mapper namespace="com.ruoyi.data.mapper.TblEquipmentMqttMapper">
 
 
+
     <resultMap type="com.ruoyi.data.domain.TblEquipmentMqtt" id="TblEquipmentMqttResult">
     <resultMap type="com.ruoyi.data.domain.TblEquipmentMqtt" id="TblEquipmentMqttResult">
+            <result property="id" column="id"/>
+            <result property="equipmentId" column="equipment_id"/>
+            <result property="mqttId" column="mqtt_id"/>
+    </resultMap>
+
+    <resultMap type="com.ruoyi.data.domain.MqttObj" id="MqttObjResult">
         <result property="id" column="id"/>
         <result property="id" column="id"/>
         <result property="equipmentId" column="equipment_id"/>
         <result property="equipmentId" column="equipment_id"/>
-        <result property="mqttId" column="mqtt_id"/>
+        <result property="protocolName" column="protocol_name"/>
+        <result property="protocolDesc" column="protocol_desc"/>
+        <result property="protocolType" column="protocol_type"/>
+        <result property="serverAddress" column="server_address"/>
+        <result property="serverTopic" column="server_topic"/>
+        <result property="account" column="account"/>
+        <result property="password" column="password"/>
+        <result property="uuid" column="uuid"/>
+        <result property="remark" column="remark"/>
+        <result property="createBy" column="create_by"/>
+        <result property="createTime" column="create_time"/>
+        <result property="updateBy" column="update_by"/>
+        <result property="updateTime" column="update_time"/>
+        <result property="status" column="status"/>
     </resultMap>
     </resultMap>
 
 
+    <select id="selectMqttListByDeviceId" parameterType="MqttObj" resultMap="MqttObjResult">
+          select a.equipment_id as equipmentId,b.* from tbl_equipment_mqtt a left join tbl_mqtt b on a.mqtt_id = b.id
+        <where>
+            <if test="equipmentId != null "> and a.equipment_id = #{equipmentId}</if>
+        </where>
+    </select>
 
 
 </mapper>
 </mapper>

+ 18 - 0
ruoyi-system/src/main/resources/mapper/data/TblRecordMapper.xml

@@ -0,0 +1,18 @@
+<?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.ruoyi.data.mapper.TblRecordMapper">
+
+    <resultMap type="com.ruoyi.data.domain.TblRecord" id="TblRecordResult">
+        <result property="id" column="id"/>
+        <result property="json" column="json"/>
+        <result property="createTime" column="create_time"/>
+        <result property="createdBy" column="created_by"/>
+        <result property="updateTime" column="update_time"/>
+        <result property="updateBy" column="update_by"/>
+        <result property="equipmentId" column="equipment_id"/>
+    </resultMap>
+
+
+</mapper>