Forráskód Böngészése

增加todo注释;
移除netty消息的处理

lixing 3 éve
szülő
commit
f1cf6f9135

+ 152 - 0
AlarmEngineStarter/src/main/java/com/persagy/apm/energyalarmstarter/alarmengine/jms/AlarmEngineMsgHandler.java

@@ -0,0 +1,152 @@
+package com.persagy.apm.energyalarmstarter.alarmengine.jms;
+
+import com.alibaba.fastjson.JSONArray;
+import com.alibaba.fastjson.JSONObject;
+import com.google.common.collect.Lists;
+import com.persagy.apm.energyalarmstarter.alarmengine.feign.AlarmCondition;
+import com.persagy.apm.energyalarmstarter.alarmengine.feign.ObjConditionRel;
+import com.persagy.apm.energyalarmstarter.alarmengine.feign.ProjectVO;
+import com.persagy.apm.energyalarmstarter.alarmengine.feign.service.AlarmServiceImpl;
+import com.persagy.apm.energyalarmstarter.alarmengine.jms.model.DmpMessage;
+import com.persagy.apm.energyalarmstarter.alarmengine.service.AlarmRecordMsgHandler;
+import com.persagy.apm.energyalarmstarter.alarmengine.util.StringUtil;
+import io.netty.channel.Channel;
+import io.netty.channel.ChannelHandler;
+import io.netty.channel.ChannelHandlerContext;
+import io.netty.channel.ChannelInboundHandlerAdapter;
+import io.netty.channel.group.ChannelGroup;
+import io.netty.channel.group.DefaultChannelGroup;
+import io.netty.util.concurrent.GlobalEventExecutor;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.lang3.StringUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.util.CollectionUtils;
+
+import java.net.SocketAddress;
+import java.util.*;
+import java.util.concurrent.ConcurrentHashMap;
+
+/**
+ * @description: Netty报警息处理中心
+ * @author: lixing
+ * @company: Persagy Technology Co.,Ltd
+ * @since: 2020/11/30 10:31 上午
+ * @version: V1.0
+ */
+@Slf4j
+@ChannelHandler.Sharable
+public class AlarmEngineMsgHandler extends ChannelInboundHandlerAdapter {
+    @Autowired
+    private AlarmRecordMsgHandler alarmRecordMsgHandler;
+
+    /**
+     * @description: 创建报警记录,并发送报警记录id至边缘端
+     * @param: nettyMessage
+     * @return: void
+     * @exception:
+     * @author: lixing
+     * @company: Persagy Technology Co.,Ltd
+     * @since: 2020/11/30 3:33 下午
+     * @version: V1.0
+     */
+    public void createAlarmRecordAndSendRecordId(DmpMessage dmpMessage) throws Exception {
+        JSONObject data = dmpMessage.getExts();
+        data.put("userId", "system");
+        String alarmRecordId = alarmRecordMsgHandler.createAlarm(data);
+        String objId = data.getString("objId");
+        String itemCode = data.getString("itemCode");
+        // TODO: 2021/11/16  把报警id放入redis缓存
+    }
+
+    /**
+     * @description: 更新报警记录
+     * @param: message
+     * @return: void
+     * @exception:
+     * @author: lixing
+     * @company: Persagy Technology Co.,Ltd
+     * @since: 2020/11/30 4:13 下午
+     * @version: V1.0
+     */
+    public void updateAlarmRecord(DmpMessage dmpMessage) throws Exception {
+        JSONObject data = dmpMessage.getExts();
+        data.put("userId", "system");
+        alarmRecordMsgHandler.updateAlarmRecord(data);
+    }
+
+    /**
+     * @description: 报警仍在持续的处理方法
+     * @param: nettyMessage
+     * @return: void
+     * @exception:
+     * @author: lixing
+     * @company: Persagy Technology Co.,Ltd
+     * @since: 2020/12/17 4:44 下午
+     * @version: V1.0
+     */
+    public void alarmContinue(DmpMessage dmpMessage) {
+        log.info("报警仍在继续:{}", dmpMessage.toString());
+    }
+
+    /**
+     * 同步新的报警条件
+     *
+     * @param msg 报警条件消息
+     * @author lixing
+     * @version V1.0 2021/10/26 7:51 下午
+     */
+    public void syncNewCondition(DmpMessage msg) {
+        AlarmCondition alarmCondition = JSONObject.parseObject(msg.getStr1(), AlarmCondition.class);
+        // TODO: 2021/11/16  更新redis
+    }
+
+    /**
+     * 同步更新的报警条件
+     *
+     * @param msg 报警条件消息
+     * @author lixing
+     * @version V1.0 2021/10/26 7:51 下午
+     */
+    public void syncUpdatedCondition(DmpMessage msg) {
+        AlarmCondition alarmCondition = JSONObject.parseObject(msg.getStr1(), AlarmCondition.class);
+        // TODO: 2021/11/16  更新redis
+    }
+
+    /**
+     * 同步删除的报警条件
+     *
+     * @param msg 报警条件消息
+     * @author lixing
+     * @version V1.0 2021/10/26 7:51 下午
+     */
+    public void syncDeletedCondition(DmpMessage msg) {
+        AlarmCondition alarmCondition = new AlarmCondition();
+        alarmCondition.setId(msg.getStr1());
+        // TODO: 2021/11/16 更新redis
+    }
+
+    /**
+     * 同步新的设备与报警条件关联关系
+     *
+     * @param msg 关联关系消息
+     * @author lixing
+     * @version V1.0 2021/10/26 7:51 下午
+     */
+    public void syncNewObjConditionRelList(DmpMessage msg) {
+        JSONArray relList = JSONObject.parseArray(msg.getStr1());
+        // TODO: 2021/11/16 更新redis
+    }
+
+    /**
+     * 同步删除的设备与报警条件关联关系
+     *
+     * @param msg 关联关系消息
+     * @author lixing
+     * @version V1.0 2021/10/26 7:51 下午
+     */
+    public void syncDeletedObjConditionRelList(DmpMessage msg) {
+        JSONArray relList = JSONObject.parseArray(msg.getStr1());
+        // TODO: 2021/11/16 更新redis
+    }
+
+}

+ 1 - 2
AlarmEngineStarter/src/main/java/com/persagy/apm/energyalarmstarter/alarmengine/jms/JmsConfig.java

@@ -1,7 +1,6 @@
 package com.persagy.apm.energyalarmstarter.alarmengine.jms;
 
 import com.persagy.apm.energyalarmstarter.alarmengine.jms.model.DmpMessage;
-import com.persagy.apm.energyalarmstarter.alarmengine.netty.NettyAlarmMsgBaseHandler;
 import com.rabbitmq.client.Channel;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.amqp.core.*;
@@ -26,7 +25,7 @@ public class JmsConfig {
      * 子类在实际项目中创建,starter中没有实例。
      */
     @Autowired
-    private NettyAlarmMsgBaseHandler msgHandler;
+    private AlarmEngineMsgHandler msgHandler;
 
     /**
      * 交换机

+ 0 - 66
AlarmEngineStarter/src/main/java/com/persagy/apm/energyalarmstarter/alarmengine/netty/NettyAlarmMessage.java

@@ -1,66 +0,0 @@
-package com.persagy.apm.energyalarmstarter.alarmengine.netty;
-
-import com.alibaba.fastjson.JSONObject;
-import com.alibaba.fastjson.annotation.JSONField;
-import lombok.Getter;
-import lombok.NoArgsConstructor;
-import lombok.Setter;
-
-import java.util.List;
-
-/**
- * @description: netty报警消息格式定义
- * @author: lixing
- * @company: Persagy Technology Co.,Ltd
- * @since: 2020/11/30 10:52 上午
- * @version: V1.0
- */
-@Getter
-@Setter
-@NoArgsConstructor
-public class NettyAlarmMessage<T> {
-    /**
-     * 唯一标识
-     */
-    @JSONField()
-    private long streamId;
-    @JSONField()
-    private int version = 1;
-
-    @JSONField()
-    private NettyMsgTypeEnum opCode;
-
-    /**
-     * 请求来源
-     */
-    @JSONField()
-    private String source = "group";
-
-    /**
-     * 用于项目控制程序清除之前的时间和命令
-     */
-    @JSONField()
-    private String clearBeforeTimeFlag;
-
-    /**
-     * 传输内容
-     */
-    @JSONField(jsonDirect = true)
-    private List<T> content;
-
-    /**
-     * 成功标识
-     */
-    @JSONField()
-    private Boolean success;
-
-    @Override
-    public String toString() {
-        return JSONObject.toJSONString(this);
-    }
-
-    public NettyAlarmMessage(NettyMsgTypeEnum opCode, List<T> content) {
-        this.opCode = opCode;
-        this.content = content;
-    }
-}

+ 0 - 533
AlarmEngineStarter/src/main/java/com/persagy/apm/energyalarmstarter/alarmengine/netty/NettyAlarmMsgBaseHandler.java

@@ -1,533 +0,0 @@
-package com.persagy.apm.energyalarmstarter.alarmengine.netty;
-
-import com.alibaba.fastjson.JSONArray;
-import com.alibaba.fastjson.JSONObject;
-import com.google.common.collect.Lists;
-import com.persagy.apm.energyalarmstarter.alarmengine.feign.AlarmCondition;
-import com.persagy.apm.energyalarmstarter.alarmengine.feign.ObjConditionRel;
-import com.persagy.apm.energyalarmstarter.alarmengine.feign.ProjectVO;
-import com.persagy.apm.energyalarmstarter.alarmengine.feign.service.AlarmServiceImpl;
-import com.persagy.apm.energyalarmstarter.alarmengine.jms.model.DmpMessage;
-import com.persagy.apm.energyalarmstarter.alarmengine.service.NettyAlarmService;
-import com.persagy.apm.energyalarmstarter.alarmengine.util.StringUtil;
-import io.netty.channel.Channel;
-import io.netty.channel.ChannelHandler;
-import io.netty.channel.ChannelHandlerContext;
-import io.netty.channel.ChannelInboundHandlerAdapter;
-import io.netty.channel.group.ChannelGroup;
-import io.netty.channel.group.DefaultChannelGroup;
-import io.netty.util.concurrent.GlobalEventExecutor;
-import lombok.extern.slf4j.Slf4j;
-import org.apache.commons.lang3.StringUtils;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.util.CollectionUtils;
-
-import java.net.SocketAddress;
-import java.util.*;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.stream.Collectors;
-
-/**
- * @description: Netty报警息处理中心
- * @author: lixing
- * @company: Persagy Technology Co.,Ltd
- * @since: 2020/11/30 10:31 上午
- * @version: V1.0
- */
-@Slf4j
-@ChannelHandler.Sharable
-public class NettyAlarmMsgBaseHandler extends ChannelInboundHandlerAdapter {
-
-    public static final String allProjects = "allProjects";
-
-    @Autowired
-    private NettyAlarmService nettyAlarmService;
-    /**
-     * 装每个客户端的地址及对应的管道
-     */
-    public Map<String, Channel> socketChannelMap = new ConcurrentHashMap<>();
-
-    @Autowired
-    private AlarmServiceImpl energyAlarmService;
-
-    /**
-     * @description: 根据项目id获取对应的通信通道
-     * @param: projectId
-     * @return: io.netty.channel.Channel
-     * @exception:
-     * @author: lixing
-     * @company: Persagy Technology Co.,Ltd
-     * @since: 2020/12/3 3:03 下午
-     * @version: V1.0
-     */
-    private Channel getChannel(String projectId) {
-        Channel channel;
-        if (StringUtils.isBlank(projectId)) {
-            channel = socketChannelMap.get(allProjects);
-            return channel;
-        }
-        // 项目上的消息只推送给一个边缘端来处理
-        channel = socketChannelMap.get(projectId);
-        if (channel == null) {
-            channel = socketChannelMap.get(allProjects);
-        }
-        return channel;
-    }
-
-    /**
-     * 保留所有与服务器建立连接的channel对象
-     */
-    public static ChannelGroup channelGroup = new DefaultChannelGroup(GlobalEventExecutor.INSTANCE);
-
-    @Override
-    public void channelRegistered(ChannelHandlerContext ctx) throws Exception {
-        log.info("netty通道注册完成,ChannelHandlerContext信息:{}", ctx.toString());
-        SocketAddress socketAddress = ctx.channel().remoteAddress();
-        String remoteAddress = socketAddress.toString();
-        System.out.println("--某个客户端绑定地址:" + remoteAddress + "--");
-    }
-
-    /**
-     * @description: 边缘端请求连接
-     * @param: nettyMessage
-     * @param: channelHandlerContext
-     * @return: void
-     * @exception:
-     * @author: lixing
-     * @company: Persagy Technology Co.,Ltd
-     * @since: 2020/11/30 11:06 上午
-     * @version: V1.0
-     */
-    private void connected(NettyAlarmMessage nettyMessage, ChannelHandlerContext channelHandlerContext) {
-        String source = nettyMessage.getSource();
-        if (StringUtils.isEmpty(source)) {
-            if (socketChannelMap.size() > 0) {
-                throw new RuntimeException("已经有projectId!=0的边缘端连接到云端,本次连接失效");
-            }
-            socketChannelMap.put(allProjects, channelHandlerContext.channel());
-        } else {
-            if (socketChannelMap.get(allProjects) != null) {
-                throw new RuntimeException("已经有projectId=0的边缘端连接到云端,本次连接失效");
-            }
-            String[] projectIds = source.split(",");
-            // 一个项目只能对应一个channel
-            for (String projectId : projectIds) {
-                // 保留旧的通道,因为当新通道请求连接时,可能老通道已经产生了通信消息
-                if (!socketChannelMap.containsKey(projectId)) {
-                    socketChannelMap.put(projectId, channelHandlerContext.channel());
-                }
-            }
-        }
-    }
-
-//    /**
-//     * @description: 发送全部报警定义到边缘端
-//     * @param: nettyMessage
-//     * @return: void
-//     * @exception:
-//     * @author: lixing
-//     * @company: Persagy Technology Co.,Ltd
-//     * @since: 2020/11/30 3:26 下午
-//     * @version: V1.0
-//     */
-//    @Deprecated
-//    public void sendAllAlarmConfigs(NettyAlarmMessage nettyMessage) throws Exception {
-//        List<JSONObject> dataList = nettyMessage.getContent();
-//        JSONObject data = dataList.get(0);
-//        String projectId = data.getString("projectId");
-//        if (StringUtils.isEmpty(projectId)) {
-//            data.put("projectId", 0);
-//        } else {
-//            data.put("projectId", projectId.split(","));
-//        }
-//
-//        data.put("userId", "system");
-//
-//        JSONArray alarmConfigs = nettyAlarmService.queryAlarmConfig(data);
-//        if (alarmConfigs == null || alarmConfigs.size() <= 0) {
-//            return;
-//        }
-//        // 查询到的报警定义按项目id分组,发送到对应的边缘端
-//        Map<String, List<Object>> groups = alarmConfigs.stream().collect(
-//                Collectors.groupingBy(alarmConfig ->
-//                        ((JSONObject) alarmConfig).getString("projectId")
-//                )
-//        );
-//
-//         /*将项目上的报警定义同步给边缘端
-//         有的边缘端可能配置了多个项目,需要将这些项目的定义先合并再发送
-//         因为边缘端如果重复接收到全量报警定义,只有最后一条生效。
-//         */
-//        Map<Channel, List<Object>> channelGroups = new HashMap<>();
-//
-//        groups.forEach((tmpProjectId, alarmConfigList) -> {
-//            Channel channel = getChannel(tmpProjectId);
-//            if (channelGroups.containsKey(channel)) {
-//                // 把要发送的报警定义追加到通道中
-//                channelGroups.get(channel).addAll(alarmConfigList);
-//            } else {
-//                channelGroups.put(channel, alarmConfigList);
-//            }
-//        });
-//
-//        channelGroups.forEach((channel, alarmConfigList) -> {
-//            sendMessage(channel, new NettyAlarmMessage(9, alarmConfigList).toString());
-//        });
-//    }
-
-    public void sendAllAlarmConditions(NettyAlarmMessage nettyMessage) throws Exception {
-        List<JSONObject> dataList = nettyMessage.getContent();
-        JSONObject data = dataList.get(0);
-        String groupCode = data.getString("groupCode");
-        String userId = "system";
-
-        Channel channel = getChannel(null);
-        // 查询报警条件,项目报警规则,项目报警规则与设备的关联关系
-        List<AlarmCondition> conditions = energyAlarmService.queryAllAlarmCondition(userId, groupCode);
-        sendMessage(channel, new NettyAlarmMessage(NettyMsgTypeEnum.ALL_CONDITIONS, conditions).toString());
-        // 报警引擎不再配置项目信息,每个报警引擎都获取全量的报警条件
-        List<ProjectVO> projects = energyAlarmService.queryProjects(userId, groupCode);
-        // 变量项目,发送每个项目上设备与报警条件的关联关系
-        for (ProjectVO project : projects) {
-            List<ObjConditionRel> objConditionRels = energyAlarmService.queryObjAlarmConditionRel(
-                    userId, groupCode, project.getProjectId());
-            if (CollectionUtils.isEmpty(objConditionRels)) {
-                continue;
-            }
-            sendMessage(channel, new NettyAlarmMessage(NettyMsgTypeEnum.ALL_OBJ_CONDITION_REL, objConditionRels).toString());
-        }
-    }
-
-    /**
-     * @description: 创建报警记录,并发送报警记录id至边缘端
-     * @param: nettyMessage
-     * @return: void
-     * @exception:
-     * @author: lixing
-     * @company: Persagy Technology Co.,Ltd
-     * @since: 2020/11/30 3:33 下午
-     * @version: V1.0
-     */
-    public void createAlarmRecordAndSendRecordId(NettyAlarmMessage nettyMessage) throws Exception {
-        List<JSONObject> dataList = nettyMessage.getContent();
-        JSONObject data = dataList.get(0);
-        data.put("userId", "system");
-        String alarmRecordId = nettyAlarmService.createAlarm(data);
-
-        JSONObject record = new JSONObject();
-        record.put("id", alarmRecordId);
-        record.put("objId", data.getString("objId"));
-        record.put("itemCode", data.getString("itemCode"));
-        List<JSONObject> records = new ArrayList<>();
-        records.add(record);
-        String projectId = data.getString("projectId");
-        sendMessage(projectId, new NettyAlarmMessage(NettyMsgTypeEnum.RECORD_ID, records).toString());
-    }
-
-    /**
-     * @description: 更新报警记录
-     * @param: message
-     * @return: void
-     * @exception:
-     * @author: lixing
-     * @company: Persagy Technology Co.,Ltd
-     * @since: 2020/11/30 4:13 下午
-     * @version: V1.0
-     */
-    public void updateAlarmRecord(NettyAlarmMessage message) throws Exception {
-        List<JSONObject> dataList = message.getContent();
-        JSONObject data = dataList.get(0);
-        data.put("userId", "system");
-        nettyAlarmService.updateAlarmRecord(data);
-    }
-
-//    /**
-//     * @description: 增量同步报警定义
-//     * @param: dmpMessage
-//     * @return: void
-//     * @exception:
-//     * @author: lixing
-//     * @company: Persagy Technology Co.,Ltd
-//     * @since: 2020/12/1 11:45 上午
-//     * @version: V1.0
-//     */
-//    public void incrementSyncAlarmConfig(DmpMessage dmpMessage) throws Exception {
-//        Map<String, JSONArray> changedAlarmConfigs = nettyAlarmService.queryChangedAlarmConfigs(dmpMessage);
-//        JSONArray createdConfigUniques = changedAlarmConfigs.get("createdConfigUniques");
-//        JSONArray deletedConfigUniques = changedAlarmConfigs.get("deletedConfigUniques");
-//        if (!CollectionUtils.isEmpty(deletedConfigUniques)) {
-//            // 通过netty发送给边缘端 10-云端推送删除的报警定义给边缘端(增量删除报警定义)
-//            sendMessage(dmpMessage.getProjectId(), new NettyAlarmMessage(10, deletedConfigUniques).toString());
-//        }
-//        if (!CollectionUtils.isEmpty(createdConfigUniques)) {
-//            // 通过netty发送给边缘端 7-云端推送修改的报警定义给边缘端(增量新增修改报警定义)
-//            sendMessage(dmpMessage.getProjectId(), new NettyAlarmMessage(7, createdConfigUniques).toString());
-//        }
-//    }
-
-    /**
-     * @param channelHandlerContext
-     * @param msg
-     * @description: 接受客户端收到的数据
-     * @return: void
-     * @exception:
-     * @author: shiliqiang
-     * @company: Persagy Technology Co.,Ltd
-     * @since: 2020/10/21 10:20
-     * @version: V1.0
-     */
-    @Override
-    public void channelRead(ChannelHandlerContext channelHandlerContext, Object msg) throws Exception {
-        log.info("收到[" + channelHandlerContext.channel().remoteAddress() + "]消息:" + msg);
-        try {
-            NettyAlarmMessage nettyMessage = StringUtil.transferItemToDTO(msg.toString(), NettyAlarmMessage.class);
-            switch (nettyMessage.getOpCode()) {
-                case CONNECT:
-                    /* 接收到边缘端创建连接请求,存储连接请求的projectId和channel映射关系,
-                     * 以后向边缘端发送请求时,通过projectId获取到对应的channel
-                     */
-                    connected(nettyMessage, channelHandlerContext);
-                    break;
-                case REQUEST_ALL_CONFIGS:
-                    // 向边缘端全量发送报警定义
-//                    sendAllAlarmConfigs(nettyMessage);
-                    sendAllAlarmConditions(nettyMessage);
-                    break;
-                case CREATE_RECORD:
-                    // 创建报警记录并发送报警记录id到边缘端
-                    createAlarmRecordAndSendRecordId(nettyMessage);
-                    break;
-                case UPDATE_RECORD_STATE:
-                    // 更新报警记录状态
-                    updateAlarmRecord(nettyMessage);
-                    break;
-                case ALARM_CONTINUE:
-                    // 报警仍在继续
-                    alarmContinue(nettyMessage);
-                    break;
-                case HEART_BEAT:
-                    // 返回一个心跳信息
-                    sendHeartBeatMsg(nettyMessage);
-                    break;
-                default:
-                    log.info("边缘端发来的参数无效,参数值为:" + nettyMessage.getOpCode());
-                    break;
-            }
-
-        } catch (Exception e) {
-            log.error("channelRead error", e);
-        }
-    }
-
-    /**
-     * 发送一个心跳消息
-     *
-     * @param nettyMessage 接收到的心跳消息
-     * @author lixing
-     * @version V1.0 2021/10/28 7:23 下午
-     */
-    private void sendHeartBeatMsg(NettyAlarmMessage nettyMessage) {
-        if (nettyMessage == null) {
-            return;
-        }
-        String projectId = nettyMessage.getSource();
-
-        sendMessage(projectId, nettyMessage.toString());
-    }
-
-    /**
-     * @description: 报警仍在持续的处理方法
-     * @param: nettyMessage
-     * @return: void
-     * @exception:
-     * @author: lixing
-     * @company: Persagy Technology Co.,Ltd
-     * @since: 2020/12/17 4:44 下午
-     * @version: V1.0
-     */
-    public void alarmContinue(NettyAlarmMessage nettyMessage) {
-        log.info("报警仍在继续:{}", nettyMessage.toString());
-    }
-
-    @Override
-    public void channelReadComplete(ChannelHandlerContext ctx) {
-        ctx.flush();
-    }
-
-
-    /**
-     * 在读取操作期间,有异常抛出时会调用。
-     *
-     * @param ctx
-     * @param cause
-     * @throws Exception
-     */
-    @Override
-    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
-        log.error(cause.getMessage());
-        cause.printStackTrace();
-        ctx.close();
-    }
-
-    /**
-     * 新增连接
-     */
-    @Override
-    public void handlerAdded(ChannelHandlerContext ctx) throws Exception {
-        // TODO Auto-generated method stub
-        //NettyServer.socketChannelMap.put(ctx.channel().remoteAddress().toString(), ctx.channel());
-        channelGroup.add(ctx.channel());
-        log.info("当前连接数:[{}],新建立连接为[{}]...", channelGroup.size(), ctx.channel().remoteAddress().toString());
-        super.handlerAdded(ctx);
-    }
-
-    /**
-     * 断开连接
-     */
-    @Override
-    public void handlerRemoved(ChannelHandlerContext ctx) throws Exception {
-        //chanel可以理解成Connection
-        Channel channel = ctx.channel();
-        Set<Map.Entry<String, Channel>> channelList = socketChannelMap.entrySet();
-        for (Map.Entry<String, Channel> channelEntry : channelList) {
-            if (channelEntry.getValue() == channel) {
-                try {
-                    // 手动关闭通道,避免异常导致的连接断开,通道未关闭
-                    if (channel.isActive()) {
-                        channel.close();
-                    }
-                } catch (Exception e) {
-                    e.printStackTrace();
-                }
-                socketChannelMap.remove(channelEntry.getKey());
-                log.warn("----项目ID[{}],地址[{}] ----离开---", channelEntry.getKey(), channel.remoteAddress().toString());
-            }
-        }
-        log.warn("----客户端[{}] ----离开", channel.remoteAddress().toString());
-    }
-
-    /**
-     * 该方法只会在通道建立时调用一次,连接生效
-     */
-    @Override
-    public void channelActive(ChannelHandlerContext ctx) throws Exception {
-        Channel channel = ctx.channel();
-
-    }
-
-    /**
-     * 连接是否有效
-     */
-    @Override
-    public void channelInactive(ChannelHandlerContext ctx) throws Exception {
-
-    }
-
-    /**
-     * @param msg
-     * @Title: sendMessage
-     * @Description: 服务端给所有客户端发送消息
-     */
-    public void sendMessageToAll(Object msg) {
-        channelGroup.writeAndFlush(msg.toString());
-    }
-
-    /**
-     * @param msg
-     * @Title: sendMessage
-     * @Description: 服务端给某个客户端发送消息
-     */
-    public void sendMessage(String projectId, String msg) {
-        Channel channel = getChannel(projectId);
-        log.info("projectId: {}", projectId);
-        sendMessage(channel, msg);
-    }
-
-    /**
-     * @description: 发送消息
-     * @param: channel 通道
-     * @param: msg 消息
-     * @return: void
-     * @exception:
-     * @author: lixing
-     * @company: Persagy Technology Co.,Ltd
-     * @since: 2020/12/17 12:55 下午
-     * @version: V1.0
-     */
-    public void sendMessage(Channel channel, String msg) {
-        if (channel != null) {
-            channel.writeAndFlush(msg);
-        } else {
-            log.error("消息通道未建立,无法发送消息!");
-        }
-    }
-
-    /**
-     * 同步新的报警条件
-     *
-     * @param msg 报警条件消息
-     * @author lixing
-     * @version V1.0 2021/10/26 7:51 下午
-     */
-    public void syncNewCondition(DmpMessage msg) {
-        AlarmCondition alarmCondition = JSONObject.parseObject(msg.getStr1(), AlarmCondition.class);
-//        JSONObject alarmCondition = JSONObject.parseObject(msg.getStr1());
-        sendMessage(msg.getProjectId(), new NettyAlarmMessage(
-                NettyMsgTypeEnum.NEW_CONDITION, Lists.newArrayList(alarmCondition)).toString());
-    }
-
-    /**
-     * 同步更新的报警条件
-     *
-     * @param msg 报警条件消息
-     * @author lixing
-     * @version V1.0 2021/10/26 7:51 下午
-     */
-    public void syncUpdatedCondition(DmpMessage msg) {
-        AlarmCondition alarmCondition = JSONObject.parseObject(msg.getStr1(), AlarmCondition.class);
-//        JSONObject alarmCondition = JSONObject.parseObject(msg.getStr1());
-        sendMessage(msg.getProjectId(), new NettyAlarmMessage(
-                NettyMsgTypeEnum.UPDATE_CONDITION, Lists.newArrayList(alarmCondition)).toString());
-    }
-
-    /**
-     * 同步删除的报警条件
-     *
-     * @param msg 报警条件消息
-     * @author lixing
-     * @version V1.0 2021/10/26 7:51 下午
-     */
-    public void syncDeletedCondition(DmpMessage msg) {
-        AlarmCondition alarmCondition = new AlarmCondition();
-        alarmCondition.setId(msg.getStr1());
-        sendMessage(msg.getProjectId(), new NettyAlarmMessage(
-                NettyMsgTypeEnum.DELETE_CONDITION, Lists.newArrayList(alarmCondition)).toString());
-    }
-
-    /**
-     * 同步新的设备与报警条件关联关系
-     *
-     * @param msg 关联关系消息
-     * @author lixing
-     * @version V1.0 2021/10/26 7:51 下午
-     */
-    public void syncNewObjConditionRelList(DmpMessage msg) {
-        JSONArray relList = JSONObject.parseArray(msg.getStr1());
-        sendMessage(msg.getProjectId(), new NettyAlarmMessage(
-                NettyMsgTypeEnum.NEW_OBJ_CONDITION_REL, Lists.newArrayList(relList)).toString());
-    }
-
-    /**
-     * 同步删除的设备与报警条件关联关系
-     *
-     * @param msg 关联关系消息
-     * @author lixing
-     * @version V1.0 2021/10/26 7:51 下午
-     */
-    public void syncDeletedObjConditionRelList(DmpMessage msg) {
-        JSONArray relList = JSONObject.parseArray(msg.getStr1());
-        sendMessage(msg.getProjectId(), new NettyAlarmMessage(
-                NettyMsgTypeEnum.DELETE_OBJ_CONDITION_REL, Lists.newArrayList(relList)).toString());
-    }
-
-}

+ 0 - 106
AlarmEngineStarter/src/main/java/com/persagy/apm/energyalarmstarter/alarmengine/netty/NettyAlarmServer.java

@@ -1,106 +0,0 @@
-package com.persagy.apm.energyalarmstarter.alarmengine.netty;
-
-import io.netty.bootstrap.ServerBootstrap;
-import io.netty.channel.ChannelFuture;
-import io.netty.channel.ChannelInitializer;
-import io.netty.channel.EventLoopGroup;
-import io.netty.channel.nio.NioEventLoopGroup;
-import io.netty.channel.socket.SocketChannel;
-import io.netty.channel.socket.nio.NioChannelOption;
-import io.netty.channel.socket.nio.NioServerSocketChannel;
-import io.netty.handler.codec.LengthFieldBasedFrameDecoder;
-import io.netty.handler.codec.LengthFieldPrepender;
-import io.netty.handler.codec.string.StringDecoder;
-import io.netty.handler.codec.string.StringEncoder;
-import io.netty.handler.logging.LogLevel;
-import io.netty.handler.logging.LoggingHandler;
-import io.netty.handler.timeout.IdleStateHandler;
-import io.netty.util.concurrent.DefaultThreadFactory;
-import io.netty.util.concurrent.OrderedEventExecutor;
-import io.netty.util.concurrent.UnorderedThreadPoolEventExecutor;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.beans.factory.annotation.Value;
-import org.springframework.stereotype.Component;
-
-import java.net.InetSocketAddress;
-import java.util.concurrent.TimeUnit;
-
-/**
- * @description: netty报警服务端
- * @author: lixing
- * @company: Persagy Technology Co.,Ltd
- * @since: 2020/11/30 10:54 上午
- * @version: V1.0
- */
-@Component
-@Slf4j
-public class NettyAlarmServer {
-    @Autowired
-    private NettyAlarmMsgBaseHandler nettyAlarmMsgBaseHandler;
-
-    @Value("${group.alarm.port}")
-    public int port;
-
-    /**
-     * 用于接收客户端的TCP连接
-     */
-    EventLoopGroup bossGroup = null;
-    /**
-     * 处理I/O相关的读写操作,或者执行系统Task、定时任务Task等。
-     */
-    EventLoopGroup workGroup = null;
-    ChannelFuture channelFuture = null;
-
-    /**
-     * @Title: start
-     * @Description: 启动netty服务端
-     */
-    public void start() {
-        log.info("NettyServer开始初始化...");
-        bossGroup = new NioEventLoopGroup();
-        workGroup = new NioEventLoopGroup();
-//        UnorderedThreadPoolEventExecutor businessGroup = new UnorderedThreadPoolEventExecutor(10, new DefaultThreadFactory("business"));
-        EventLoopGroup businessGroup = new NioEventLoopGroup();
-        try {
-            // (服务端启动类)ServerBootstrap负责初始化netty服务器,并且开始监听端口的socket请求
-            ServerBootstrap startNetty = new ServerBootstrap();
-            startNetty.group(bossGroup, workGroup)
-                    .channel(NioServerSocketChannel.class)
-                    .option(NioChannelOption.SO_BACKLOG, 1024)
-                    .childOption(NioChannelOption.TCP_NODELAY, true)
-                    .handler(new LoggingHandler(LogLevel.INFO))
-                    .childHandler(new ChannelInitializer<SocketChannel>() {
-                        @Override
-                        protected void initChannel(SocketChannel ch) throws Exception {
-                            // 解决粘包和拆包问题,增加的编码和解码器
-                            ch.pipeline().addLast(new LengthFieldBasedFrameDecoder(Integer.MAX_VALUE, 0, 4, 0, 4));
-                            ch.pipeline().addLast(new LengthFieldPrepender(4));
-                            // the encoder and decoder are static as these are sharable
-                            ch.pipeline().addLast(new StringDecoder());
-                            ch.pipeline().addLast(new StringEncoder());
-                            // 为监听客户端read/write事件的Channel添加用户自定义的ChannelHandler
-                            ch.pipeline().addLast(businessGroup, nettyAlarmMsgBaseHandler);
-
-                            // 增加超时检查
-                            ch.pipeline().addLast(new IdleStateHandler(10, 0, 0, TimeUnit.SECONDS));
-                        }
-                    });
-            channelFuture = startNetty.bind(new InetSocketAddress(port)).sync();
-
-            log.info("NettyServer初始化完成,启动netty服务端");
-            channelFuture.channel().closeFuture().sync();
-
-        } catch (Exception e) {
-            e.printStackTrace();
-        } finally {
-            // 关闭主线程组
-            bossGroup.shutdownGracefully();
-            // 关闭工作线程组
-            workGroup.shutdownGracefully();
-            businessGroup.shutdownGracefully();
-        }
-    }
-
-
-}

+ 0 - 41
AlarmEngineStarter/src/main/java/com/persagy/apm/energyalarmstarter/alarmengine/netty/NettyMsgTypeEnum.java

@@ -1,41 +0,0 @@
-package com.persagy.apm.energyalarmstarter.alarmengine.netty;
-
-import lombok.AllArgsConstructor;
-import lombok.Getter;
-import lombok.Setter;
-
-/**
- * netty的消息类型
- *
- * @author lixing
- * @version V1.0 2021/10/25 3:30 下午
- */
-@AllArgsConstructor
-public enum NettyMsgTypeEnum {
-
-    /**
-     * netty的消息类型
-     */
-    HEART_BEAT(10, "心跳包"),
-    ACCEPTED(100, "已接收到消息"),
-    CONNECT(200, "建立连接,此时的source == 项目id"),
-    REQUEST_ALL_CONFIGS(10, "边缘端申请全量获取报警定义(报警条件、条件和设备的关联关系)"),
-    ALL_CONDITIONS(11, "全量报警条件"),
-    NEW_CONDITION(12, "新增报警条件"),
-    UPDATE_CONDITION(13, "更新报警条件"),
-    DELETE_CONDITION(14, "删除报警条件"),
-    ALL_OBJ_CONDITION_REL(21, "全量条件和设备的关联关系"),
-    NEW_OBJ_CONDITION_REL(22, "新增条件和设备的关联关系"),
-    DELETE_OBJ_CONDITION_REL(23, "删除条件和设备的关联关系"),
-    CREATE_RECORD(31, "创建报警记录"),
-    RECORD_ID(32, "报警记录id"),
-    UPDATE_RECORD_STATE(33, "更新报警记录状态"),
-    ALARM_CONTINUE(34, "报警仍在持续");
-
-    @Setter
-    @Getter
-    private int value;
-    @Setter
-    @Getter
-    private String desc;
-}

+ 0 - 31
AlarmEngineStarter/src/main/java/com/persagy/apm/energyalarmstarter/alarmengine/netty/runner/NettyServerRunner.java

@@ -1,31 +0,0 @@
-package com.persagy.apm.energyalarmstarter.alarmengine.netty.runner;
-
-import com.persagy.apm.energyalarmstarter.alarmengine.netty.NettyAlarmServer;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.boot.CommandLineRunner;
-import org.springframework.core.annotation.Order;
-import org.springframework.stereotype.Component;
-
-
-/**
- * @description: 初始化项目
- * @author: shiliqiang
- * @company: Persagy Technology Co.,Ltd
- * @since: 2020/10/20 15:37
- * @version: V1.0
- */
-@Component
-@Slf4j
-@Order(Integer.MAX_VALUE)
-public class NettyServerRunner implements CommandLineRunner {
-	@Autowired
-	private NettyAlarmServer nettyServer;
-
-	@Override
-	public void run(String... args) throws Exception {
-		// 项目加载完成后启动netty
-		log.info("-------启动NettyAlarmServer--------");
-		nettyServer.start();
-	}
-}

+ 6 - 7
AlarmEngineStarter/src/main/java/com/persagy/apm/energyalarmstarter/alarmengine/service/NettyAlarmService.java

@@ -19,14 +19,13 @@ import java.util.Map;
 import java.util.stream.Collectors;
 
 /**
- * @description: 处理netty消息中调用数据中台的逻辑
- * @author: lixing
- * @company: Persagy Technology Co.,Ltd
- * @since: 2020/11/30 2:38 下午
- * @version: V1.0
- **/
+ * 处理报警引擎发送的报警创建、更新、持续消息
+ *
+ * @author lixing
+ * @version V1.0 2021/11/16 8:24 下午
+ */
 @Slf4j
-public abstract class NettyAlarmService extends BaseService {
+public abstract class AlarmRecordMsgHandler extends BaseService {
     @Autowired
     AlarmClient alarmClient;