zhongrj
2023-07-06 8453121fb84156efe06d313aadc3281013f0b359
新增 sse 消息推送
1 files modified
3 files added
208 ■■■■■ changed files
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/constant/ZtApiUrlConstant.java 2 ●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/controller/SSEController.java 41 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/server/SSEServer.java 148 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/vo/SseVO.java 17 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/constant/ZtApiUrlConstant.java
@@ -38,7 +38,7 @@
    public static final String over_z_res_api = "/services/1234567890ABCDEFGHIJKLMN/over_z_res";
    /**
     *
     * 水位超设计洪水位告警处理
     */
    public static final String over_des_rz_warn_api = "/services/1234567890ABCDEFGHIJKLMN/over_des_rz/warn";
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/controller/SSEController.java
New file
@@ -0,0 +1,41 @@
package cn.gistack.alerts.sse.controller;
import cn.gistack.alerts.sse.server.SSEServer;
import cn.gistack.alerts.sse.vo.SseVO;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
@Slf4j
@RestController
@CrossOrigin
@RequestMapping("/sse/sse")
public class SSEController {
    /**
     * 建立连接
     * @param sse
     * @return
     */
    @GetMapping("/connect")
    public SseEmitter connect(SseVO sse){
        return SSEServer.connect(sse.getType() + ":" + sse.getUserId());
    }
    /**
     * 发送消息
     * @throws InterruptedException
     */
    @GetMapping("/process")
    public void sendMessage() throws InterruptedException {
        SSEServer.sendMessage("web:123","hello!");
//        for(int i=0; i<=100; i++){
//            if(i>50&&i<70){
//                Thread.sleep(500L);
//            }else{
//                Thread.sleep(100L);
//            }
//            SSEServer.batchSendMessage(String.valueOf(i));
//        }
    }
}
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/server/SSEServer.java
New file
@@ -0,0 +1,148 @@
package cn.gistack.alerts.sse.server;
import lombok.extern.slf4j.Slf4j;
import org.springframework.http.MediaType;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;
@Slf4j
public class SSEServer {
    /**
     * 当前连接数
     */
    private static AtomicInteger count = new AtomicInteger(0);
    private static Map<String, SseEmitter> sseEmitterMap = new ConcurrentHashMap<>();
    public static SseEmitter connect(String userId){
        //设置超时时间,0表示不过期,默认是30秒,超过时间未完成会抛出异常
        SseEmitter sseEmitter = new SseEmitter(0L);
        //注册回调
        sseEmitter.onCompletion(completionCallBack(userId));
        sseEmitter.onError(errorCallBack(userId));
        sseEmitter.onTimeout(timeOutCallBack(userId));
        sseEmitterMap.put(userId,sseEmitter);
        //数量+1
        count.getAndIncrement();
        log.info("create new sse connect ,current user:{}",userId);
        return sseEmitter;
    }
    /**
     * 给指定用户发消息
     */
    public static void sendMessage(String userId, String message){
        if(sseEmitterMap.containsKey(userId)){
            try{
                sseEmitterMap.get(userId).send(message);
            }catch (IOException e){
                log.error("user id:{}, send message error:{}",userId,e.getMessage());
                e.printStackTrace();
            }
        }
    }
    /**
     * 想多人发送消息,组播
     */
    public static void groupSendMessage(String groupId, String message){
        if(sseEmitterMap!=null&&!sseEmitterMap.isEmpty()){
            sseEmitterMap.forEach((k,v) -> {
                try{
                    if(k.startsWith(groupId)){
                        v.send(message, MediaType.APPLICATION_JSON);
                    }
                }catch (IOException e){
                    log.error("user id:{}, send message error:{}",groupId,message);
                    removeUser(k);
                }
            });
        }
    }
    /**
     * 批量发送消息
     * @param message
     */
    public static void batchSendMessage(String message) {
        sseEmitterMap.forEach((k,v)->{
            try{
                v.send(message,MediaType.APPLICATION_JSON);
            }catch (IOException e){
                log.error("user id:{}, send message error:{}",k,e.getMessage());
                removeUser(k);
            }
        });
    }
    /**
     * 群发消息
     */
    public static void batchSendMessage(String message, Set<String> userIds){
        userIds.forEach(userId->sendMessage(userId,message));
    }
    /**
     * 用户离线删除用户
     * @param userId
     */
    public static void removeUser(String userId){
        sseEmitterMap.remove(userId);
        //数量-1
        count.getAndDecrement();
        log.info("remove user id:{}",userId);
    }
    public static List<String> getIds(){
        return new ArrayList<>(sseEmitterMap.keySet());
    }
    public static int getUserCount(){
        return count.intValue();
    }
    /**
     * 结束回调
     * @param userId
     * @return
     */
    private static Runnable completionCallBack(String userId) {
        return () -> {
            log.info("结束连接,{}",userId);
            removeUser(userId);
        };
    }
    /**
     * 超时回调
     * @param userId
     * @return
     */
    private static Runnable timeOutCallBack(String userId){
        return ()->{
            log.info("连接超时,{}",userId);
            removeUser(userId);
        };
    }
    /**
     * 错误回调
     * @param userId
     * @return
     */
    private static Consumer<Throwable> errorCallBack(String userId){
        return throwable -> {
            log.error("连接异常,{}",userId);
            removeUser(userId);
        };
    }
}
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/vo/SseVO.java
New file
@@ -0,0 +1,17 @@
package cn.gistack.alerts.sse.vo;
import lombok.Data;
@Data
public class SseVO {
    /**
     * 类型 web,app,小程序
     */
    private String type;
    /**
     * 用户唯一值
     */
    private String userId;
}