服务端推送技术 Server-sent Events springBoot代码示例

简介: 服务端推送技术 Server-sent Events springBoot代码示例

 SSE推送技术

SSE全称Server-sent Events,是HTML 5 规范的一个组成部分,具体去MDN网站查看相关文档。该规范十分简单,

SSE推送技术是服务器端与浏览器端之间的通讯协议,通讯协议是基于纯文本的简单协议。

服务器端的响应的内容类型是“text/event-stream”。响应文本的内容可以看成是一个事件流,由不同的事件所组成。每个事件由类型和数据两部分组成,同时每个事件可以有一个可选的标识符。不同事件的内容之间通过仅包含回车符和换行符的空行(“rn”)来分隔。每个事件的数据可能由多行组成。

image.gif编辑

如上图所示,每个事件之间通过空行来分隔。每一行都是由键值对组成。如果键为空则表示该行为注释,会在处理时被忽略。例如第10行。第1行表示一个只包含数据的事件。会按照默认事件走(message事件)。第3-4代表一个附带eventID的事件。第6-8代表一个自定义事件。第10-14代表一个多行数据事件,多行数据由换行符链接

key定义有以下几种:

    • data,表示该行包含的是数据。以 data 开头的行可以出现多次。所有这些行都是该事件的数据。
    • 类型为 event,表示该行用来声明事件的类型。浏览器在收到数据时,会产生对应类型的事件。默认提供三个标准事件(当然你可以自定义):

    image.gif编辑

      • id,表示该行用来声明事件的标识符。服务器端返回的数据中包含了事件的标识符,浏览器会记录最近一次接收到的事件的标识符。如果与服务器端的连接中断,当浏览器端再次进行连接时,会通过 HTTP 头“Last-Event-ID”来声明最后一次接收到的事件的标识符。服务器端可以通过浏览器端发送的事件标识符来确定从哪个事件开始来继续连接。
      • retry,表示该行用来声明浏览器在连接断开之后进行再次连接之前的等待时间。

      SSE只适用于高级浏览器,但是注意IE不直接支持。IE上的XMLHttpRequest对象不支持获取部分的响应内容,所以不支持。每次总有IE怪不得快被淘汰了。

      SSE VS Websocket

        • SSE 只能Server到Client单项,而Websocket是双向通信。
        • SSE 比 Websocket 轻量。当然功能要简单的多。开发便利,不牵涉协议升级问题。
        • SSE 天然支持断线重连

        服务端代码示例

        import com.alibaba.fastjson.JSONObject;
        import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
        import com.baomidou.mybatisplus.core.toolkit.CollectionUtils;
        import com.hxtx.spacedata.common.domain.ResponseDTO;
        import com.hxtx.spacedata.domain.entity.task.TaskInfoEntity;
        import com.hxtx.spacedata.enums.task.TaskInfoStatusEnum;
        import com.hxtx.spacedata.mapper.task.TaskInfoDao;
        import lombok.extern.slf4j.Slf4j;
        import org.springframework.beans.factory.annotation.Autowired;
        import org.springframework.scheduling.annotation.Scheduled;
        import org.springframework.web.bind.annotation.GetMapping;
        import org.springframework.web.bind.annotation.PathVariable;
        import org.springframework.web.bind.annotation.RestController;
        import javax.servlet.http.HttpServletRequest;
        import javax.servlet.http.HttpSession;
        import java.util.Iterator;
        import java.util.List;
        import java.util.Map;
        import java.util.concurrent.ConcurrentHashMap;
        import java.util.stream.Collectors;
        /**
         * 服务端推送技术 server-sent events
         * @description
         * @author tarzan Liu
         * @version 1.0.0
         * @date 2020/10/27
         */
        @RestController
        @Slf4j
        public class SSEController {
            @Autowired
            private TaskInfoDao taskInfoDao;
            private static ConcurrentHashMap<String,Long> ssePushUsers = new ConcurrentHashMap<>();
            /**
             *  如果没有客户端,则直接修改消息已发送 (2分钟执行一次)
             * @author sunboqiang
             * @date 2020/11/3
             */
            @Scheduled(cron = "0 0/2 * * * ?")
            public void finishSend() {
                if(ssePushUsers.size()==0){
                    QueryWrapper<TaskInfoEntity> queryWrapper = new QueryWrapper<>();
                    queryWrapper.lambda().eq(TaskInfoEntity::getStatus, TaskInfoStatusEnum.SUCCESS.getStatus());
                    queryWrapper.lambda().eq(TaskInfoEntity::getSendStatus,0);
                    List<TaskInfoEntity> list = taskInfoDao.selectList(queryWrapper);
                    if(CollectionUtils.isNotEmpty(list)){
                        taskInfoDao.updateSendStatusByIds(list.stream().map(TaskInfoEntity::getId).collect(Collectors.toList()), 2);
                    }
                }
            }
            /**
             *  剔除关闭的客户端
             * @author sunboqiang
             * @date 2020/11/3
             */
            @Scheduled(cron = "0/2 * * * * ?") // 2S执行一次
            public void clear() {
                //2秒执行一次,时间差>5S 说明客户端关闭了,直接剔除
                long now = System.currentTimeMillis();
                for (Iterator<Map.Entry<String, Long>> it = ssePushUsers.entrySet().iterator(); it.hasNext(); ) {
                    Map.Entry<String, Long> item = it.next();
                    long time = item.getValue();
                    //log.info(item.getKey()+"注册时间差:"+(now - time)/1000);
                    if(now - time > 5000){
                        //5 秒
                        it.remove();
                        log.info("剔除客户端:"+item.getKey());
                    }
                }
            }
            @GetMapping(value="/sse/push/version/get")
            public String getVersion(HttpServletRequest request){
                HttpSession session = request.getSession();
                if(null != session){
                    return session.getId();
                }
                return null;
            }
            /**
             *  推送C++ json文件编译情况信息
             * @author sunboqiang
             * @date 2020/10/29
             */
            @GetMapping(value="/sse/push/{version}",produces="text/event-stream;charset=utf-8")
            public String push(@PathVariable("version") String version) {
                try {
                    Thread.sleep(1000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                QueryWrapper<TaskInfoEntity> queryWrapper = new QueryWrapper<>();
                queryWrapper.lambda().eq(TaskInfoEntity::getStatus, TaskInfoStatusEnum.SUCCESS.getStatus());
                queryWrapper.lambda().eq(TaskInfoEntity::getSendStatus,0);
                List<TaskInfoEntity> list = taskInfoDao.selectList(queryWrapper);
                String data = "";
                if(CollectionUtils.isEmpty(list)){
                    //还没有消息,收集等待推送的客户端
                    ssePushUsers.put(version,System.currentTimeMillis());
                    //data = "data:没有编译消息,当前打开客户端数量:"+ ssePushUsers.size()+"个;" +"\n\n";
                } else {
                    List<Long> drawingIds = list.stream().map(TaskInfoEntity::getDrawingId).distinct().collect(Collectors.toList());
                    //编译成功,推送消息
                    if(ssePushUsers.size()>0){
                        //存在接收客户端
                        ResponseDTO result = new ResponseDTO();
                        result.setCode(1);
                        result.setMsg("有新的编译");
                        result.setSuccess(true);
                        result.setData("drawingIds="+drawingIds);
                        data = "data:"+ JSONObject.toJSONString(result) +"\n\n";
                        ssePushUsers.remove(version);
                        if(ssePushUsers.size() == 0){
                            //最后一个客户端推送完成
                            taskInfoDao.updateSendStatusByIds(list.stream().map(TaskInfoEntity::getId).collect(Collectors.toList()), 1);
                        }
                    } else {
                        //没有客户端,直接推送成功
                        taskInfoDao.updateSendStatusByIds(list.stream().map(TaskInfoEntity::getId).collect(Collectors.toList()), 1);
                    }
                }
                return data;
            }
        }

        image.gif


        相关文章
        |
        安全 Java 应用服务中间件
        Spring Boot + Java 21:内存减少 60%,启动速度提高 30% — 零代码
        通过调整三个JVM和Spring Boot配置开关,无需重写代码即可显著优化Java应用性能:内存减少60%,启动速度提升30%。适用于所有在JVM上运行API的生产团队,低成本实现高效能。
        1315 3
        |
        网络协议 Java
        SpringBoot快速搭建TCP服务端和客户端
        由于工作需要,研究了SpringBoot搭建TCP通信的过程,对于工程需要的小伙伴,只是想快速搭建一个可用的服务.其他的教程看了许多,感觉讲得太复杂,很容易弄乱,这里我只讲效率,展示快速搭建过程。
        1321 58
        |
        监控 Java 数据安全/隐私保护
        阿里面试:SpringBoot启动时, 如何执行扩展代码?你们项目 SpringBoot 进行过 哪些 扩展?
        阿里面试:SpringBoot启动时, 如何执行扩展代码?你们项目 SpringBoot 进行过 哪些 扩展?
        |
        前端开发 Java 物联网
        智慧班牌源码,采用Java + Spring Boot后端框架,搭配Vue2前端技术,支持SaaS云部署
        智慧班牌系统是一款基于信息化与物联网技术的校园管理工具,集成电子屏显示、人脸识别及数据交互功能,实现班级信息展示、智能考勤与家校互通。系统采用Java + Spring Boot后端框架,搭配Vue2前端技术,支持SaaS云部署与私有化定制。核心功能涵盖信息发布、考勤管理、教务处理及数据分析,助力校园文化建设与教学优化。其综合性和可扩展性有效打破数据孤岛,提升交互体验并降低管理成本,适用于日常教学、考试管理和应急场景,为智慧校园建设提供全面解决方案。
        842 70
        |
        Java 数据库连接 数据库
        Spring boot 使用mybatis generator 自动生成代码插件
        本文介绍了在Spring Boot项目中使用MyBatis Generator插件自动生成代码的详细步骤。首先创建一个新的Spring Boot项目,接着引入MyBatis Generator插件并配置`pom.xml`文件。然后删除默认的`application.properties`文件,创建`application.yml`进行相关配置,如设置Mapper路径和实体类包名。重点在于配置`generatorConfig.xml`文件,包括数据库驱动、连接信息、生成模型、映射文件及DAO的包名和位置。最后通过IDE配置运行插件生成代码,并在主类添加`@MapperScan`注解完成整合
        1783 1
        Spring boot 使用mybatis generator 自动生成代码插件
        |
        Java 数据库连接 API
        Java 8 + 特性及 Spring Boot 与 Hibernate 等最新技术的实操内容详解
        本内容涵盖Java 8+核心语法、Spring Boot与Hibernate实操,按考试考点分类整理,含技术详解与代码示例,助力掌握最新Java技术与应用。
        424 2
        SpringBoot快速搭建WebSocket服务端和客户端
        由于工作需要,研究了SpringBoot搭建WebSocket双向通信的过程,其他的教程看了许多,感觉讲得太复杂,很容易弄乱,这里我只展示快速搭建过程。
        3079 1
        |
        Java 调度 流计算
        基于Java 17 + Spring Boot 3.2 + Flink 1.18的智慧实验室管理系统核心代码
        这是一套基于Java 17、Spring Boot 3.2和Flink 1.18开发的智慧实验室管理系统核心代码。系统涵盖多协议设备接入(支持OPC UA、MQTT等12种工业协议)、实时异常检测(Flink流处理引擎实现设备状态监控)、强化学习调度(Q-Learning算法优化资源分配)、三维可视化(JavaFX与WebGL渲染实验室空间)、微服务架构(Spring Cloud构建分布式体系)及数据湖建设(Spark构建实验室数据仓库)。实际应用中,该系统显著提升了设备调度效率(响应时间从46分钟降至9秒)、设备利用率(从41%提升至89%),并大幅减少实验准备时间和维护成本。
        671 0
        |
        XML 前端开发 Java
        SpringBoot整合Flowable【04】- 通过代码控制流程流转
        本文介绍了如何使用Flowable的Java API控制流程流转,基于前文构建的绩效流程模型。首先,通过Flowable-UI导出模型文件并部署到Spring Boot项目中。接着,详细讲解了如何通过代码部署、启动和审批流程,涉及`RepositoryService`、`RuntimeService`和`TaskService`等核心服务类的使用。最后,通过实际操作演示了流程从部署到完成的全过程,并简要说明了相关数据库表的变化。本文帮助读者初步掌握Flowable在实际业务中的应用,后续将深入探讨更多高级功能。
        2627 0
        SpringBoot整合Flowable【04】-  通过代码控制流程流转
        |
        JavaScript 安全 Java
        java版药品不良反应智能监测系统源码,采用SpringBoot、Vue、MySQL技术开发
        基于B/S架构,采用Java、SpringBoot、Vue、MySQL等技术自主研发的ADR智能监测系统,适用于三甲医院,支持二次开发。该系统能自动监测全院患者药物不良反应,通过移动端和PC端实时反馈,提升用药安全。系统涵盖规则管理、监测报告、系统管理三大模块,确保精准、高效地处理ADR事件。
        721 1