Quellcode durchsuchen

新增实时及登录推送站内信

lyh vor 1 Jahr
Ursprung
Commit
6bdd9f2aa5

+ 4 - 1
src/main/java/com/ydtech/SpringInsApplication.java

@@ -1,9 +1,11 @@
 package com.ydtech;
 
+import com.ydtech.modules.admin.controller.websocket.SendMessageServer;
 import org.mybatis.spring.annotation.MapperScan;
 import org.springframework.boot.SpringApplication;
 import org.springframework.boot.autoconfigure.SpringBootApplication;
 import org.springframework.cache.annotation.EnableCaching;
+import org.springframework.context.ConfigurableApplicationContext;
 import org.springframework.scheduling.annotation.EnableAsync;
 import org.springframework.scheduling.annotation.EnableScheduling;
 
@@ -37,7 +39,8 @@ import org.springframework.scheduling.annotation.EnableScheduling;
 public class SpringInsApplication {
 
     public static void main(String[] args) {
-        SpringApplication.run(SpringInsApplication.class, args);
+        ConfigurableApplicationContext run = SpringApplication.run(SpringInsApplication.class, args);
+        SendMessageServer.setApplicationContext(run);
     }
 
 }

+ 1 - 1
src/main/java/com/ydtech/modules/admin/constants/MessageModelConstants.java

@@ -9,7 +9,7 @@ public enum MessageModelConstants {
 
     //渠道管理/出单员管理
     //%添加人员账号姓名%于%创建时间%新增了出单员,%出单员姓名%,%出单员工号%,请及时查看。
-    MODEL_1("1", "有新增的出单员", "${添加人员账号姓名}于${创建时间}新增了出单员,${出单员姓名},${出单员工号},请及时查看。", "/peopleManage/player/playerExamine"),
+    MODEL_1("1", "有新增的出单员", "${添加人员账号姓名}于${创建时间}新增了出单员,${出单员姓名},${出单员工号},请及时查看。", "/peopleManage/player/playerManage"),
     //代理人管理/代理人审核
     //%人员类型% %待审核人员姓名%,%待审核人员工号%,于%审核数据生成时间%申请入司,请及时审批。
     MODEL_2("2", "有新的入司人员,请注意查收", "${人员类型} ${待审核人员姓名},${待审核人员工号}于${审核数据生成时间}申请入司,请及时审批。", "/peopleManage/agent/agentExamine");

+ 194 - 0
src/main/java/com/ydtech/modules/admin/controller/websocket/SendMessageServer.java

@@ -0,0 +1,194 @@
+package com.ydtech.modules.admin.controller.websocket;
+
+import com.github.pagehelper.PageInfo;
+import com.google.gson.Gson;
+import com.google.gson.JsonParser;
+import com.ydtech.modules.message.entity.SystemMessage;
+import com.ydtech.modules.message.entity.SystemMessageVo;
+import com.ydtech.modules.message.entity.dto.SystemMessageDTO;
+import com.ydtech.modules.message.service.SystemMessageService;
+import com.ydtech.utils.StringUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.context.ApplicationContext;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import javax.websocket.*;
+import javax.websocket.server.PathParam;
+import javax.websocket.server.ServerEndpoint;
+import com.alibaba.fastjson.JSON;
+import com.alibaba.fastjson.JSONObject;
+import cn.hutool.log.Log;
+import cn.hutool.log.LogFactory;
+
+
+@ServerEndpoint(value = "/getMessage/{userId}")
+@Component
+public class SendMessageServer {
+
+
+    private static ApplicationContext applicationContext;
+
+    public static void setApplicationContext(ApplicationContext context) {
+        applicationContext = context;
+    }
+
+    static Log log = LogFactory.get(SendMessageServer.class);
+    /**静态变量,用来记录当前在线连接数。应该把它设计成线程安全的。*/
+    private static int onlineCount = 0;
+    /**concurrent包的线程安全Set,用来存放每个客户端对应的MyWebSocket对象。*/
+    private static ConcurrentHashMap<String,SendMessageServer> webSocketMap = new ConcurrentHashMap<>();
+    /**与某个客户端的连接会话,需要通过它来给客户端发送数据*/
+    private Session session;
+    /**接收userId*/
+    private String userId="";
+
+    /**
+     * 连接建立成功调用的方法*/
+    @OnOpen
+    public void onOpen(Session session, @PathParam("userId") String userId) {
+        this.session = session;
+        this.userId=userId;
+        if(webSocketMap.containsKey(userId)){
+            webSocketMap.remove(userId);
+            webSocketMap.put(userId, this);
+            //加入set中
+        }else{
+            webSocketMap.put(userId, this);
+            //加入set中
+            addOnlineCount();
+            //在线数加1
+        }
+
+        log.info("用户连接:"+userId+",当前在线人数为:" + getOnlineCount());
+
+        try {
+            PageInfo<SystemMessageDTO> message = getMessage(userId);
+            log.info("为用户" + userId + "连接建立成功发送站内信" + new Gson().toJson(message));
+            sendMessage(new Gson().toJson(message));
+        } catch (Exception e) {
+            log.error("用户:"+userId+",网络异常!!!!!!");
+        }
+    }
+
+    /**
+     * 连接关闭调用的方法
+     */
+    @OnClose
+    public void onClose() {
+        if(webSocketMap.containsKey(userId)){
+            webSocketMap.remove(userId);
+            //从set中删除
+            subOnlineCount();
+        }
+        log.info("用户退出:"+userId+",当前在线人数为:" + getOnlineCount());
+    }
+
+    /**
+     * 收到客户端消息后调用的方法
+     *
+     * @param message 客户端发送过来的消息*/
+    @OnMessage
+    public void onMessage(String message, Session session) {
+        log.info("用户消息:"+userId+",报文:"+message);
+        //可以群发消息
+        //消息保存到数据库、redis
+        if (StringUtils.isNotBlank(message)) {
+            try {
+                //解析发送的报文
+                JSONObject jsonObject = JSON.parseObject(message);
+                //追加发送人(防止串改)
+                jsonObject.put("fromUserId",this.userId);
+                String toUserId=jsonObject.getString("userId");
+
+                //传送给对应toUserId用户的websocket
+                if (StringUtils.isNotBlank(toUserId) && webSocketMap.containsKey(toUserId)) {
+                    PageInfo<SystemMessageDTO> getMessage = getMessage(toUserId);
+                    log.info("为用户" + toUserId + "登录或刷新发送站内信" + new Gson().toJson(getMessage));
+                    webSocketMap.get(toUserId).sendMessage(new Gson().toJson(getMessage));
+                } else {
+                    log.error("请求的userId:"+toUserId+"不在该服务器上");
+                    //否则不在这个服务器上,发送到mysql或者redis
+                }
+            } catch (Exception e) {
+                e.printStackTrace();
+            }
+        }
+    }
+
+    /**
+     *
+     * @param session
+     * @param error
+     */
+    @OnError
+    public void onError(Session session, Throwable error) {
+        log.error("用户错误:"+this.userId+",原因:"+error.getMessage());
+        error.printStackTrace();
+    }
+    /**
+     * 实现服务器主动推送
+     */
+    public void sendMessage(String message) throws Exception {
+        this.session.getBasicRemote().sendText(message);
+    }
+
+
+    /**
+     * 每小时发送站内信消息
+     * */
+//    @Scheduled(cron = "0 * * * * *")
+//    public void fixedTime() throws Exception {
+//        log.info("定时任务发送站内信消息到:" + JSONObject.toJSONString(webSocketMap.keySet()));
+//        if (webSocketMap.size() > 0) {
+//            for (String userId : webSocketMap.keySet()) {
+//                PageInfo<SystemMessageDTO> message = getMessage(userId);
+//                log.info("定时任务发送站内信消息内容:" + JSONObject.toJSONString(message));
+//                webSocketMap.get(userId).sendMessage(new Gson().toJson(message));
+//            }
+//        }
+//    }
+
+
+    public void realTime(List<SystemMessage> systemMessageList) throws Exception {
+        log.info("实时发送站内信消息到:" + JSONObject.toJSONString(webSocketMap.keySet()));
+        if (webSocketMap.size() > 0) {
+            for (String userId : webSocketMap.keySet()) {
+                for (SystemMessage systemMessage : systemMessageList) {
+                    if (userId.equals(systemMessage.getReceiverId())) {
+                        log.info("实时发送站内信消息内容:" + JSONObject.toJSONString(systemMessage));
+                        webSocketMap.get(userId).sendMessage(new Gson().toJson(systemMessage));
+                    }
+                }
+            }
+        }
+    }
+
+
+    public static synchronized int getOnlineCount() {
+        return onlineCount;
+    }
+
+    public static synchronized void addOnlineCount() {
+        SendMessageServer.onlineCount++;
+    }
+
+    public static synchronized void subOnlineCount() {
+        SendMessageServer.onlineCount--;
+    }
+
+    public PageInfo<SystemMessageDTO> getMessage(String userId) {
+        SystemMessageVo systemMessageVo = new SystemMessageVo();
+        systemMessageVo.setUserId(userId);
+        systemMessageVo.setPageNum(1);
+        systemMessageVo.setPageSize(10);
+        SystemMessageService socketTableConnService = applicationContext.getBean(SystemMessageService.class);
+        PageInfo<SystemMessageDTO> systemMessageByUser = socketTableConnService.getSystemMessageByUser(systemMessageVo);
+        return systemMessageByUser;
+    }
+
+
+}

+ 6 - 0
src/main/java/com/ydtech/modules/message/service/impl/SystemMessageServiceImpl.java

@@ -6,6 +6,7 @@ import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
 import com.github.pagehelper.PageHelper;
 import com.github.pagehelper.PageInfo;
 import com.ydtech.modules.admin.constants.MessageConstants;
+import com.ydtech.modules.admin.controller.websocket.SendMessageServer;
 import com.ydtech.modules.message.dao.SystemMessageMapper;
 import com.ydtech.modules.message.entity.SystemMessage;
 import com.ydtech.modules.message.entity.SystemMessageVo;
@@ -15,6 +16,7 @@ import com.ydtech.utils.StringUtils;
 import lombok.extern.slf4j.Slf4j;
 import org.apache.commons.lang3.ObjectUtils;
 import org.springframework.beans.BeanUtils;
+import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.stereotype.Service;
 
 import java.util.ArrayList;
@@ -27,6 +29,9 @@ import java.util.regex.Pattern;
 @Slf4j
 public class SystemMessageServiceImpl extends ServiceImpl<SystemMessageMapper, SystemMessage> implements SystemMessageService {
 
+    @Autowired
+    private SendMessageServer sendMessageServer;
+
     @Override
     public PageInfo<SystemMessageDTO> getSystemMessageByUser(SystemMessageVo systemMessageVo){
         log.info("获取站内信入参:{}", JSONObject.toJSONString(systemMessageVo));
@@ -64,6 +69,7 @@ public class SystemMessageServiceImpl extends ServiceImpl<SystemMessageMapper, S
         try {
             if (ObjectUtils.isNotEmpty(systemMessageList) && !systemMessageList.isEmpty()) {
                 saveBatch(systemMessageList);
+                sendMessageServer.realTime(systemMessageList);
             }
         } catch (Exception e) {
             e.printStackTrace();