WebSocket
快来分享你的内容吧~
- 2024-03-10·后端开发前端 WebSocket 的一些使用 WebSocket 是一种网络通信协议,用于实现双向通信。在前端中,你可以使用 JavaScript 中的 WebSocket 对象来创建 WebSocket 连接,发送和接收数据。 连接的建立 通过创建一个 WebSocket 对象建立一个 WebSocket 连接 例如: const ws = new WebSocket('ws://localhost:8查看全文鱼友0412:本文非常优质!欢迎鱼友参与知识碎片的贡献,并获取奖励:https://docs.qq.com/form/page/DQkdHTGFLQmJjV3VU#/fill知识碎片汇总:https://yuyuanweb.feishu.cn/wiki/AqfawFUT0iD69kkiRKoci6Nqnqc1910分享
- 2023-08-29·前端开发
WebSocket实现聊天
## 一、先补个基础:WebSocket 是什么? 在传统的 HTTP 通信中,客户端(比如浏览器)要获取新消息,只能 “主动问” 服务器(比如每隔 30 秒发一次请求查新消息),这种方式实时性差、浪费资源。 而 **WebSocket 是一种 “全双工” 通信协议**:客户端和服务器只需建立一次连接,之后双方可以随时互相发消息(就像打电话,接通后能随时说话),完美适配 “私信实时收发” 的场景。 这段代码就是 **服务器端的 WebSocket 服务实现**,负责管理用户连接、接收用户发送的私信、实时推送给接收者。 ## 二、代码结构拆解:每个部分是干嘛的? 先看整体结构:代码用 `@ServerEndpoint` 注解标记这是一个 WebSocket 服务,包含「连接管理」「消息处理」「异常处理」「主动推送」四大核心功能,还有一些静态变量用于存储状态。 ### 1. 静态变量:存储全局状态 ```java // 1. 线程池:处理“异步存库”(避免存数据库的耗时操作阻塞实时推送) private static final ExecutorService EXECUTOR_SERVICE = Executors.newFixedThreadPool(10); // 2. 存储用户连接:key=用户ID,value=WebSocket会话(Session) // 用 ConcurrentHashMap 是因为多用户并发操作时安全(避免线程问题) private static final Map<Long, Session> USER_SESSIONS = new ConcurrentHashMap<>(); // 3. 服务层注入:WebSocket 是多实例的,所以用静态变量存储 Service(否则注入会失败) private static PrivateMessageService privateMessageService; private static UserService userService; // 4. 静态注入的 setter 方法:Spring 会调用这个方法给静态 Service 赋值 @Resource public void setPrivateMessageService(PrivateMessageService service) { PrivateMessageWebSocket.privateMessageService = service; } @Resource public void setUserService(UserService service) { PrivateMessageWebSocket.userService = service; } ``` **关键理解**: - WebSocket 实例是 “多例” 的(每个用户连接都会创建一个实例),而 Spring 的 Service 是 “单例” 的,所以不能直接用 `@Autowired` 注入,必须通过「静态变量 + 静态 setter」的方式注入 Service。 - `USER_SESSIONS` 是核心:记录当前哪些用户在线(有 WebSocket 连接),后续推送消息时要靠它找到接收者的连接。 ### 2. 核心注解方法:WebSocket 生命周期回调 WebSocket 有固定的生命周期(建立连接→收发消息→关闭连接 / 异常),代码用以下注解方法对应这些生命周期,当触发对应事件时,Spring 会自动调用这些方法。 #### (1)@OnOpen:用户建立 WebSocket 连接时触发 ```java /** * 用户打开页面、建立WebSocket连接时调用(比如进入私信聊天页) * @param session 当前用户的WebSocket会话(相当于“通话通道”,用于后续发消息) * @param userId 从URL路径中获取的用户ID(比如连接地址是 ws://xxx/ws/private-message/100,这里userId就是100) */ @OnOpen public void onOpen(Session session, @PathParam("userId") Long userId) { // 把“用户ID”和“他的会话”存入全局map,标记该用户已在线 USER_SESSIONS.put(userId, session); log.info("用户[{}]建立WebSocket连接,当前在线数:{}", userId, USER_SESSIONS.size()); } ``` **场景举例**:用户 A(ID=100)打开私信页面,前端会发起连接请求 `new WebSocket("ws://localhost:8080/ws/private-message/100")`,服务器收到后调用 `onOpen`,把 100 和对应的 `session` 存到 `USER_SESSIONS` 里,此时在线数 + 1。 #### (2)@OnMessage:接收客户端(用户)发送的私信时触发 这是最核心的方法,负责 “接收消息→存数据库→推送给接收者” 三步,代码里已经分了步骤,我们逐行拆: ```java /** * 用户发送私信时调用(比如用户A给用户B发“你好”,前端会通过WebSocket把消息发给服务器) * @param message 前端传过来的消息内容(JSON格式,对应 PrivateMessageSendDTO) * @param senderId 发送者ID(从URL路径获取,比如发送者是100) */ @OnMessage public void onMessage(String message, @PathParam("userId") Long senderId) { try { // 步骤1:解析前端传的JSON消息,转成Java对象(PrivateMessageSendDTO) // 比如前端发的JSON是 {"receiverId":200,"content":"你好","msgType":1},这里会转成DTO PrivateMessageSendDTO dto = JSON.parseObject(message, PrivateMessageSendDTO.class); Long receiverId = dto.getReceiverId(); // 接收者ID(比如200) // 步骤2:校验参数(避免空消息、没填接收者的情况) if (ObjectUtils.isEmpty(receiverId) || ObjectUtils.isEmpty(dto.getContent())) { log.warn("私信参数无效:senderId={}, message={}", senderId, message); return; } // 步骤3:异步写入数据库(重点!先推送再存库,保证实时性) // 为什么用线程池异步?因为存数据库是耗时操作(比如100ms),如果同步执行,会阻塞后续的推送,导致接收者延迟收到消息 EXECUTOR_SERVICE.submit(() -> { // 调用 PrivateMessageService 的方法,把私信存到数据库(包含标记“最新消息”、设置过期时间等逻辑) privateMessageService.sendPrivateMessage(senderId, dto); }); // 步骤4:实时推送消息给接收者(如果接收者在线) // 从 USER_SESSIONS 里查接收者(比如200)是否有在线连接 Session receiverSession = USER_SESSIONS.get(receiverId); if (receiverSession != null && receiverSession.isOpen()) { // 构建推送用的VO(给前端返回的格式,包含发送者昵称、头像等,方便前端显示) User sender = userService.getById(senderId); // 查发送者的用户信息 String senderNickname = sender != null ? sender.getNickname() : "匿名用户"; String senderAvatar = sender != null ? sender.getAvatar() : ""; PrivateMessageVO pushVO = new PrivateMessageVO(); pushVO.setSenderId(senderId); // 发送者ID pushVO.setSenderNickname(senderNickname); // 发送者昵称 pushVO.setSenderAvatar(senderAvatar); // 发送者头像 pushVO.setReceiverId(receiverId); // 接收者ID pushVO.setContent(dto.getContent()); // 消息内容 pushVO.setMsgType(dto.getMsgType()); // 消息类型(1=文本,2=图片) pushVO.setIsRead(0); // 初始未读 pushVO.setSendTime(new Date()); // 发送时间 // 把VO转成JSON,通过接收者的Session推送给前端 receiverSession.getBasicRemote().sendText(JSON.toJSONString(pushVO)); log.info("私信推送成功:senderId={}, receiverId={}", senderId, receiverId); } else { // 如果接收者不在线(USER_SESSIONS里没有他的Session),只存数据库,等他上线后再补推 log.info("接收者[{}]不在线,私信将存入数据库", receiverId); } } catch (Exception e) { // 捕获所有异常,避免单个消息处理失败导致整个WebSocket服务崩溃 log.error("处理私信消息失败:senderId={}, message={}", senderId, message, e); } } ``` **场景举例**: 用户 A(100)给用户 B(200)发 “你好”,前端把消息转成 JSON 发给服务器,服务器调用 `onMessage`: 1. 解析出接收者是 200,内容是 “你好”; 2. 用线程池异步把这条消息存到数据库; 3. 查 `USER_SESSIONS`,如果 B 在线(有 Session),就把包含 A 昵称、头像的消息推给 B 的前端,B 页面实时显示 “你好”;如果 B 不在线,只存库,等 B 下次上线再补推。 #### (3)@OnClose:用户断开 WebSocket 连接时触发 ```java /** * 用户关闭页面、断开连接时调用(比如关闭私信页、退出登录) * @param userId 断开连接的用户ID * @param session 要关闭的会话 */ @OnClose public void onClose(@PathParam("userId") Long userId, Session session) { // 从 USER_SESSIONS 中移除该用户,标记为离线 USER_SESSIONS.remove(userId); try { // 关闭会话(释放资源) session.close(); } catch (IOException e) { log.error("关闭会话失败:userId={}", userId, e); } log.info("用户[{}]断开WebSocket连接,当前在线数:{}", userId, USER_SESSIONS.size()); } ``` **场景举例**:用户 A 关闭私信页面,前端会主动断开 WebSocket 连接,服务器调用 `onClose`,把 A 的 ID 从 `USER_SESSIONS` 中移除,在线数 - 1。 #### (4)@OnError:WebSocket 连接异常时触发 ```java /** * 连接出现异常时调用(比如网络断了、前端崩溃) * @param userId 异常用户的ID * @param session 异常的会话 * @param throwable 异常信息 */ @OnError public void onError(@PathParam("userId") Long userId, Session session, Throwable throwable) { log.error("用户[{}]WebSocket连接异常", userId, throwable); // 异常时也要移除连接(避免存无效的Session) USER_SESSIONS.remove(userId); try { session.close(); // 关闭异常会话 } catch (IOException e) { log.error("异常关闭会话失败:userId={}", userId, e); } } ``` **作用**:处理意外情况(比如用户网络突然断开),避免无效的连接占用服务器资源,同时记录异常日志方便排查问题。 ### 3. 主动推送方法:pushMessage(非回调,手动调用) ```java /** * 主动推送消息给用户(比如用户重连后,补推他离线时收到的未读消息) * @param userId 要推送的用户ID * @param vo 要推送的私信VO */ public static void pushMessage(Long userId, PrivateMessageVO vo) { // 查用户是否在线 Session session = USER_SESSIONS.get(userId); if (session != null && session.isOpen()) { try { // 推送消息 session.getBasicRemote().sendText(JSON.toJSONString(vo)); log.info("主动推送消息给用户[{}]成功", userId); } catch (IOException e) { log.error("主动推送消息失败:userId={}", userId, e); } } } ``` **使用场景**:比如用户 B 之前离线,收到了 3 条未读消息,当他再次上线建立 WebSocket 连接后,服务器可以调用这个方法,把这 3 条未读消息主动推送给 B,让 B 一上线就能看到。 ## 三、完整流程串讲:从 “用户发消息” 到 “对方收消息” 用一个具体场景把所有环节串起来,你会更清楚: 假设 **用户 A(ID=100)给用户 B(ID=200)发私信 “在吗?”**,整个流程如下: 1. **建立连接**: - A 打开私信页面,前端发起 WebSocket 连接:`new WebSocket("ws://localhost:8080/ws/private-message/100")`; - 服务器调用 `@OnOpen`,把 `100 → A的Session` 存到 `USER_SESSIONS`,在线数 = 1; - (如果 B 也打开了页面)B 同样建立连接,服务器把 `200 → B的Session` 存到 `USER_SESSIONS`,在线数 = 2。 2. **A 发送消息**: - A 在前端输入 “在吗?”,点击发送; - 前端把消息转成 JSON(比如 `{"receiverId":200,"content":"在吗?","msgType":1}`),通过 A 的 WebSocket 发给服务器; - 服务器调用 `@OnMessage`,解析 JSON 得到 DTO,校验参数(接收者 200 非空、内容非空)。 3. **异步存库**: - 服务器用 `EXECUTOR_SERVICE` 线程池,异步调用 `privateMessageService.sendPrivateMessage(100, dto)`; - Service 层会做两件事:① 把之前 A 和 B 会话的 “最新消息” 标记为旧(`isLatest=0`);② 插入这条新消息(`isLatest=1`,过期时间 = 当前 + 30 天,`isRead=0`)。 4. **实时推送给 B**: - 服务器从 `USER_SESSIONS` 中查 200 的 Session(如果 B 在线,能查到); - 构建 `PrivateMessageVO`(包含 A 的昵称、头像、消息内容等),转成 JSON 推给 B 的 Session; - B 的前端收到 JSON,解析后在页面上显示 “A 发来消息:在吗?”。 5. **如果 B 不在线**: - 服务器查不到 200 的 Session,只执行 “异步存库”,不推送; - 等 B 下次上线建立连接后,服务器可以调用 `pushMessage(200, 未读消息VO)`,把离线时的消息补推给 B。 6. **A 关闭页面**: - 前端断开连接,服务器调用 `@OnClose`,从 `USER_SESSIONS` 中移除 100,在线数 = 1(只剩 B)。 ## 四、关键细节:为什么要这么设计? 1. **为什么用线程池异步存库?** 存数据库是耗时操作(比如 IO 读写需要 50-100ms),如果同步执行,会阻塞 `@OnMessage` 方法,导致推送消息延迟。用线程池异步处理,能让 “推送消息” 优先执行,保证实时性。 2. **为什么用 ConcurrentHashMap 存连接?** 多用户会并发建立 / 断开连接(比如同时有 100 个用户上线),普通 `HashMap` 在并发操作时会出现线程安全问题(比如死循环),`ConcurrentHashMap` 是线程安全的,适合这种场景。 3. **为什么 Service 要用静态注入?** WebSocket 是 “多实例” 的(每个用户连接都会 new 一个 `PrivateMessageWebSocket` 对象),而 Spring 的 Service 是 “单例” 的(整个应用只有一个 `PrivateMessageService` 实例)。如果用普通 `@Autowired` 注入,每个 WebSocket 实例都会拿到一个新的 Service 实例,这不符合 Spring 设计,还可能导致数据不一致。所以用 “静态变量 + 静态 setter” 的方式,让所有 WebSocket 实例共享同一个 Service 实例。 4. **为什么要标记 “最新消息”(isLatest=1)?** 会话列表(比如 “与 B 的对话”“与 C 的对话”)只需要显示每条对话的最新一条消息,不需要查所有历史消息。存库时标记 `isLatest=1`,后续查会话列表时只查 `isLatest=1` 的数据,能大幅优化查询性能。 ## 五、前端怎么配合? 服务器端的 WebSocket 服务需要前端配合才能工作,前端核心代码大概是这样(以 JavaScript 为例): ```javascript // 1. 建立连接(userId 是当前登录用户的ID,比如100) const userId = 100; const ws = new WebSocket(`ws://localhost:8080/ws/private-message/${userId}`); // 2. 连接成功时触发 ws.onopen = () => { console.log("WebSocket连接成功"); }; // 3. 接收服务器推送的消息(比如收到A发来的“在吗?”) ws.onmessage = (event) => { const messageVO = JSON.parse(event.data); // 解析服务器推送的VO console.log("收到消息:", messageVO); // 把消息显示到页面上(比如添加到聊天记录列表) renderMessage(messageVO); }; // 4. 发送消息(给接收者200发“在吗?”) function sendMessage() { const dto = { receiverId: 200, content: "在吗?", msgType: 1 // 文本类型 }; ws.send(JSON.stringify(dto)); // 转成JSON发给服务器 } // 5. 连接关闭时触发 ws.onclose = () => { console.log("WebSocket连接关闭"); // 可选:自动重连(比如网络断了后重试) reconnect(); }; // 6. 连接异常时触发 ws.onerror = (error) => { console.error("WebSocket异常:", error); }; ``` ## 总结 这段代码的核心是 **“用 WebSocket 实现私信的实时收发”**,通过 `@ServerEndpoint` 定义服务端点,用 `@OnOpen/@OnMessage/@OnClose/@OnError` 处理连接生命周期,用 `USER_SESSIONS` 管理在线用户,用线程池异步存库保证实时性,最终实现 “用户发消息→服务器实时推送给接收者” 的完整流程。 如果还有某个细节没懂(比如 Service 层的 `sendPrivateMessage` 逻辑、前端重连机制),可以再针对性问~
通过Netty手写Web服务器
最近通过原生netty写了一个web服务器,实现了一些在SpringBoot中常用的功能(我觉得对初学者有很好的开拓效果,能够清楚的感知市面大部分web服务器的技术细节),还需要大家多给一些意见。。。 https://gitee.com/wlx0902/netty-server.git
最快速学会websocket、spring security、swagger+knife4j、openfeign的使用
<html> <head></head> <body> <div class="content ql-editor"> <p>本篇教程直指实战,10分钟掌握websocket、spring security结合jwt形式token、swagger结合kinfe4j、openfeign与springboot的集成使用</p> <p>websocket的使用</p> <p>引入依赖</p> <div class="ql-code-block-container"> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>org.springframework.boot</groupId> </div> <div class="ql-code-block"> <artifactId>spring-boot-starter-websocket</artifactId> </div> <div class="ql-code-block"> </dependency> </div> <div class="ql-code-block"> 版本号在项目父pom.xml文件中可以看到,后续依赖同理 </div> </div> <p>注册websocket端点</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @Configuration </div> <div class="ql-code-block"> @EnableWebSocket </div> <div class="ql-code-block"> public class WebSocketConfig implements WebSocketConfigurer { </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Resource </div> <div class="ql-code-block"> WebSocketHandler webSocketHandler; </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { </div> <div class="ql-code-block"> 将处理器注册进容器,并指定连接路径为服务器ip:Springboot端口/路径(localhost:1688/test/test)的websocket连接交由处理器处理,并设置允许跨域 </div> <div class="ql-code-block"> registry.addHandler(webSocketHandler, “/test/test”).setAllowedOrigins("*"); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <p>创建处理器并注册进容器</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @Slf4j </div> <div class="ql-code-block"> @Component </div> <div class="ql-code-block"> public class WebSocketHandler extends TextWebSocketHandler { </div> <div class="ql-code-block"> 重写handleTextMessage方法,有一条websocket消息发送到Springboot服务时可以获取到该消息的websocketsession对象和封装了发送消息的textmessage对象 </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { </div> <div class="ql-code-block"> 用gson将该对象中蕴含的websocket消息相关信息如消息具体内容转换成json格式对象的属性 </div> <div class="ql-code-block"> JsonObject json = JsonParser.parseString(message.getPayload()).getAsJsonObject(); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> 重写建立连接后会调用的方法 </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public void afterConnectionEstablished(WebSocketSession session) throws Exception { </div> <div class="ql-code-block"> log.info("连接建立"); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> 重写关闭连接后会调用的方法 </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public void afterConnectionClosed(WebSocketSession session, CloseStatus status) { </div> <div class="ql-code-block"> log.info("连接关闭"); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <p>spring security的使用</p> <p>引入依赖</p> <div class="ql-code-block-container"> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>org.springframework.boot</groupId> </div> <div class="ql-code-block"> <artifactId>spring-boot-starter-security</artifactId> </div> <div class="ql-code-block"> </dependency> </div> </div> <p>创建核心功能配置类,由于gateway是webflux风格,集成security写法和普通mvc集成略有不同,下面是普通mvc集成写法</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @Configuration </div> <div class="ql-code-block"> @EnableWebSecurity </div> <div class="ql-code-block"> public class SecurityConfig extends WebSecurityConfigurerAdapter { </div> <div class="ql-code-block"> @Resource </div> <div class="ql-code-block"> JwtAuthenticationTokenFilter jwtAuthenticationTokenFilter; </div> <div class="ql-code-block"> 重写核心配置方法自定义权限规则 </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> protected void configure(HttpSecurity http) throws Exception { </div> <div class="ql-code-block"> 链式写法 </div> <div class="ql-code-block"> http </div> <div class="ql-code-block"> .authorizeRequests() .antMatchers("/getAllUser","/doc.html","/index.html","/login","/vue.js","/axios.min.js","/indexs.html","/video/getvideo").permitAll() 像doc.html、index.html这些允许匿名访问 </div> <div class="ql-code-block"> .antMatchers("/path1").hasAuthority("menu:insert")而path1、path2这些路径需要有 </div> <div class="ql-code-block"> .antMatchers("/path2").hasAuthority("menu:delete")menu:insert或menu:delete的权限 </div> <div class="ql-code-block"> .anyRequest().authenticated() 任何请求都需要接受规则验证 </div> <div class="ql-code-block"> .and() </div> <div class="ql-code-block"> .formLogin() 表单登录禁止,前后端分离时需要禁止,否则每次打开页面都会进入security </div> <div class="ql-code-block"> .disable() 自带表单 </div> <div class="ql-code-block"> .formLogin() </div> <div class="ql-code-block"> .loginPage("/index.html") 前后端不分离时定义登录页面所在位置 </div> <div class="ql-code-block"> .loginProcessingUrl("/loginTest") 识别登录的请求 </div> <div class="ql-code-block"> .successHandler(new JwtAuthenticationSuccessHandler()) 登录成功处理器 </div> <div class="ql-code-block"> .defaultSuccessUrl("/page", true) 成功后转向的页面 </div> <div class="ql-code-block"> .failureHandler(new CustomAuthenticationFailureHandler()) 失败的处理器 </div> <div class="ql-code-block"> 该方案由security封装实现,自动识别登录,无需编写登录接口,但直接编写登录接口也能起到同样作用且可以少掉重写security授权服务步骤,因此后续都以自己实现登录接口做示例 </div> <div class="ql-code-block"> .and() </div> <div class="ql-code-block"> .httpBasic().disable() 禁用http基本认证 </div> <div class="ql-code-block"> .csrf().disable() 禁用跨域伪造请求保护,因为需要使用jwt </div> <div class="ql-code-block"> .addFilterAt(jwtAuthenticationTokenFilter, UsernamePasswordAuthenticationFilter.class); 用自定义鉴权过滤器替换掉security默认鉴权过滤器 </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <p>自定义登录接口调用的service层核心方法</p> <div class="ql-code-block-container"> <div class="ql-code-block"> private String login(LoginRequest loginRequest){ </div> <div class="ql-code-block"> 查询数据库中是否存在和该用户名密码匹配的用户 </div> <div class="ql-code-block"> LambdaQueryWrapper<User> userWrapper=new LambdaQueryWrapper<>(); </div> <div class="ql-code-block"> userWrapper.eq(User::getUsername,loginRequest.getUsername()); </div> <div class="ql-code-block"> User user=userMapper.selectOne(userWrapper); </div> <div class="ql-code-block"> 如果有 </div> <div class="ql-code-block"> if(user.getPassword().equals(loginRequest.getPassword())){ </div> <div class="ql-code-block"> 查询该用户有的权限,以容易理解的形式,权限表中只存用户id和拥有的权限 </div> <div class="ql-code-block"> LambdaQueryWrapper<Authority> authorityWrapper=new LambdaQueryWrapper<>(); </div> <div class="ql-code-block"> authorityWrapper.eq(Authority::getUserId,user.getId()); </div> <div class="ql-code-block"> List<String> authorities=authorityMapper.selectList(authorityWrapper); </div> <div class="ql-code-block"> 将该用户拥有的系列权限和userid一起封装进token返回给controller,再在controller中将token保存到响应头中返回给前端 </div> <div class="ql-code-block"> return JwtUtil.generateToken(user.getId(),authorities); </div> <div class="ql-code-block"> }else{ </div> <div class="ql-code-block"> 若登录失败则返回的token为空 </div> <div class="ql-code-block"> return null; </div> <div class="ql-code-block"> } </div> </div> <p>鉴权过滤器</p> <p><br></p> <div class="ql-code-block-container"> <div class="ql-code-block"> 前端发送请求时将jwt携带到请求头里 </div> <div class="ql-code-block"> String jwt=request.getHeader("authorization"); </div> <div class="ql-code-block"> if (jwt != null) { </div> <div class="ql-code-block"> 解析JWT令牌并获取权限 </div> <div class="ql-code-block"> Authentication authentication = JwtUtil.getAuthentication(jwt); </div> <div class="ql-code-block"> 设置认证信息到SecurityContextHolder,用来后续验证调接口时是否有该接口的权限 </div> <div class="ql-code-block"> SecurityContextHolder.getContext().setAuthentication(authentication); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> chain.doFilter(request, response); </div> <div class="ql-code-block"> } </div> </div> <p>而gateway集成security如下</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @Configuration </div> <div class="ql-code-block"> @EnableWebSecurity </div> <div class="ql-code-block"> public class SecurityConfig extends WebSecurityConfigurerAdapter { </div> <div class="ql-code-block"> @Resource </div> <div class="ql-code-block"> JwtAuthorizationFilter jwtAuthorizationFilter; </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> @Bean </div> <div class="ql-code-block"> public SecurityWebFilterChain springSecurityFilterChain(ServerHttpSecurity http) { </div> <div class="ql-code-block"> http </div> <div class="ql-code-block"> .authorizeExchange() 起始授权、后续填写匹配路径和鉴定权限方法与mvc风格不一致 </div> <div class="ql-code-block"> .pathMatchers( "/webjars/","/register/**","/user-center/**", </div> <div class="ql-code-block"> "/swagger-ui.html", </div> <div class="ql-code-block"> "/webjars/**", </div> <div class="ql-code-block"> "/swagger-resources/**", </div> <div class="ql-code-block"> "/v2/**") </div> <div class="ql-code-block"> .permitAll() </div> <div class="ql-code-block"> .anyExchange().hasAuthority("role:user") </div> <div class="ql-code-block"> .and() </div> <div class="ql-code-block"> .httpBasic().disable() </div> <div class="ql-code-block"> .formLogin().disable() </div> <div class="ql-code-block"> .addFilterAt(new JwtAuthorizationFilter(),SecurityWebFiltersOrder.AUTHORIZATION) </div> <div class="ql-code-block"> .csrf().disable(); </div> <div class="ql-code-block"> return http.build(); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <p>而gateway集成security的自定义鉴权过滤器如下</p> <div class="ql-code-block-container"> <div class="ql-code-block"> public class JwtAuthorizationFilter implements WebFilter { </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> /** </div> <div class="ql-code-block"> *打印请求路径,用自定义请求头隔绝csrf攻击,取出token认证用户与验证权限 </div> <div class="ql-code-block"> */ </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) { </div> <div class="ql-code-block"> String path = exchange.getRequest().getURI().getPath(); </div> <div class="ql-code-block"> if (!path.contains("webjars") && !path.contains("swagger") && !path.contains("api")) { </div> <div class="ql-code-block"> log.info(exchange.getRequest().getURI().getPath()); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> 从请求头中获取token </div> <div class="ql-code-block"> String jwt = exchange.getRequest().getHeaders().getFirst(Constant.SHORT_TOKEN); </div> <div class="ql-code-block"> if (jwt != null) { </div> <div class="ql-code-block"> 获取权限 </div> <div class="ql-code-block"> Claims claims = JwtUtil.getClaimsFromToken(jwt); </div> <div class="ql-code-block"> String role = claims.get(JWT_ROLE, String.class); </div> <div class="ql-code-block"> SimpleGrantedAuthority authority = new SimpleGrantedAuthority("role"+role); </div> <div class="ql-code-block"> Authentication authentication = new UsernamePasswordAuthenticationToken( </div> <div class="ql-code-block"> "username", </div> <div class="ql-code-block"> "password", </div> <div class="ql-code-block"> Collections.singletonList(authority) </div> <div class="ql-code-block"> ); </div> <div class="ql-code-block"> SecurityContext context = new SecurityContextImpl(); </div> <div class="ql-code-block"> 将权限封装成authentication对象后放入安全上下文 </div> <div class="ql-code-block"> context.setAuthentication(authentication); </div> <div class="ql-code-block"> return chain.filter(exchange) .contextWrite(ReactiveSecurityContextHolder.withSecurityContext(Mono.just(context))); </div> </div> <p>jwt的使用</p> <p>引入依赖</p> <div class="ql-code-block-container"> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>io.jsonwebtoken</groupId> </div> <div class="ql-code-block"> <artifactId>jjwt</artifactId> </div> <div class="ql-code-block"> </dependency> </div> </div> <p>创建工具类</p> <div class="ql-code-block-container"> <div class="ql-code-block"> package example.util; </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> import io.jsonwebtoken.*; </div> <div class="ql-code-block"> import org.springframework.security.authentication.UsernamePasswordAuthenticationToken; </div> <div class="ql-code-block"> import org.springframework.security.core.authority.SimpleGrantedAuthority; </div> <div class="ql-code-block"> import org.springframework.security.core.Authentication; </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> import java.util.Date; </div> <div class="ql-code-block"> import java.util.List; </div> <div class="ql-code-block"> import java.util.stream.Collectors; </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> public class JwtUtil { </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> private static final String SECRET_KEY = "key"; </div> <div class="ql-code-block"> public static String generateToken(Integer userId, List<SimpleGrantedAuthority> authorities) { </div> <div class="ql-code-block"> List<String> authoritiesList = authorities.stream() </div> <div class="ql-code-block"> .map(SimpleGrantedAuthority::getAuthority) </div> <div class="ql-code-block"> .collect(Collectors.toList()); </div> <div class="ql-code-block"> int expirationTime=10000; </div> <div class="ql-code-block"> 使用建造者模式一步步构建jwt </div> <div class="ql-code-block"> return Jwts.builder() </div> <div class="ql-code-block"> .setSubject(String.valueOf(userId)) 设置主题 </div> <div class="ql-code-block"> .claim("authorities", authoritiesList) 设置声明,也是jwt体中重要信息的存储位置 </div> <div class="ql-code-block"> .setIssuedAt(new Date(System.currentTimeMillis())) 设置jwt创建时间 </div> <div class="ql-code-block"> .setExpiration(new Date(System.currentTimeMillis()+expirationTime)) 设置jwt过期时间 </div> <div class="ql-code-block"> .signWith(SignatureAlgorithm.HS512, SECRET_KEY) 设置签名规则和签名key,需要根据key来解码jwt </div> <div class="ql-code-block"> .compact(); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> public static Authentication getAuthentication(String token) { </div> <div class="ql-code-block"> 去掉jwt中的前缀 </div> <div class="ql-code-block"> token = token.replace("Bearer ", ""); </div> <div class="ql-code-block"> 根据预先存储的key解码jwt获取jwt的核心信息如所具备权限 </div> <div class="ql-code-block"> Claims claims = Jwts.parser() </div> <div class="ql-code-block"> .setSigningKey(SECRET_KEY) </div> <div class="ql-code-block"> .parseClaimsJws(token) </div> <div class="ql-code-block"> .getBody(); </div> <div class="ql-code-block"> List<String> authoritiesList = claims.get("authorities", List.class); </div> <div class="ql-code-block"> List<SimpleGrantedAuthority> authorities = authoritiesList.stream() </div> <div class="ql-code-block"> .map(SimpleGrantedAuthority::new) </div> <div class="ql-code-block"> .collect(Collectors.toList()); </div> <div class="ql-code-block"> String username = claims.getSubject(); </div> <div class="ql-code-block"> 将权限设置入上下文 </div> <div class="ql-code-block"> return new UsernamePasswordAuthenticationToken(username, null, authorities); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> public static boolean validateToken(String token) { </div> <div class="ql-code-block"> try { </div> <div class="ql-code-block"> 解析Token的同时会检查签名。如果签名错误,会抛出SignatureException </div> <div class="ql-code-block"> 如果Token已过期,会抛出ExpiredJwtException </div> <div class="ql-code-block"> Jwts.parser() </div> <div class="ql-code-block"> .setSigningKey(SECRET_KEY) </div> <div class="ql-code-block"> .parseClaimsJws(token); </div> <div class="ql-code-block"> 如果没有抛出异常,那么Token就是既有效又没有过期的 </div> <div class="ql-code-block"> return true; </div> <div class="ql-code-block"> } catch (ExpiredJwtException | SignatureException e) { </div> <div class="ql-code-block"> Token过期或签名错误 </div> <div class="ql-code-block"> return false; </div> <div class="ql-code-block"> } catch (JwtException e) { </div> <div class="ql-code-block"> 其他可能的错误,例如构造的JWT格式不正确等 </div> <div class="ql-code-block"> return false; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <p>最后是security封装好的加密工具使用</p> <div class="ql-code-block-container"> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> 配置类中注册该工具类 </div> <div class="ql-code-block"> public PasswordEncoder passwordEncoder() { </div> <div class="ql-code-block"> return new BCryptPasswordEncoder(); </div> <div class="ql-code-block"> } </div> </div> <p>加密工具使用(默认使用加盐算法)</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @Service </div> <div class="ql-code-block"> public class UserService { </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Autowired </div> <div class="ql-code-block"> private PasswordEncoder passwordEncoder; </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> public void registerUser(String rawPassword) { </div> <div class="ql-code-block"> 将密码加密后保存到数据库 </div> <div class="ql-code-block"> String encodedPassword = passwordEncoder.encode(rawPassword); </div> <div class="ql-code-block"> saveUserToDatabase(encodedPassword); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> 实现将用户信息保存到数据库 </div> <div class="ql-code-block"> private void saveUserToDatabase(String encodedPassword) { </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> 验证请求中的密码和数据库中保存的加密密码解密后的值是否一致 </div> <div class="ql-code-block"> public boolean authenticate(String rawPassword, String storedEncodedPassword) { </div> <div class="ql-code-block"> return passwordEncoder.matches(rawPassword, storedEncodedPassword); </div> <div class="ql-code-block"> } </div> </div> <p>swagger文档结合knife4j的使用</p> <p>引入依赖</p> <div class="ql-code-block-container"> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>com.github.xiaoymin</groupId> </div> <div class="ql-code-block"> <artifactId>knife4j-openapi2-spring-boot-starter</artifactId> </div> <div class="ql-code-block"> </dependency> </div> </div> <p>核心配置类</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @Configuration </div> <div class="ql-code-block"> @EnableSwagger2WebMvc </div> <div class="ql-code-block"> public class Knife4jConfiguration { </div> <div class="ql-code-block"> @Bean(value = "dockerBean") </div> <div class="ql-code-block"> public Docket dockerBean() { </div> <div class="ql-code-block"> Docket docket=new Docket(DocumentationType.SWAGGER_2) </div> <div class="ql-code-block"> .apiInfo(new ApiInfoBuilder() </div> <div class="ql-code-block"> .title("视频模块文档") 接口文档标题以及一些相关描述信息 </div> <div class="ql-code-block"> .description(" 首页、视频详情页、个人主页视频数据的获取以及对视频的操作") </div> <div class="ql-code-block"> .termsOfServiceUrl("https://doc.Amumu.com/") </div> <div class="ql-code-block"> .contact("阿沐木") </div> <div class="ql-code-block"> .version("1.0") </div> <div class="ql-code-block"> .build()) </div> <div class="ql-code-block"> .groupName("视频服务") </div> <div class="ql-code-block"> .select() ↓核心配置,标注接口文档的接口映射路径 </div> <div class="ql-code-block"> .apis(RequestHandlerSelectors.basePackage("ljl.bilibili.video.controller")) </div> <div class="ql-code-block"> .paths(PathSelectors.any()) </div> <div class="ql-code-block"> .build(); </div> <div class="ql-code-block"> return docket; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <p>接口文档映射路径在idea中可以找到</p> <p><img src="https://pic.code-nav.cn/planet_post_image/1808693273146032129/0ba6p12m.jpeg"></p> <p>找到要放入接口文档的接口所在的包点击,可以在上方看到路径,src-main-java后的路径就是映射路径(tips:点击红色圆圈所在位置可以快速找到一个文件在项目中的具体位置)</p> <p>在请求类中描述信息</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @Data </div> <div class="ql-code-block"> 描述该类 </div> <div class="ql-code-block"> @ApiModel("收藏请求") </div> <div class="ql-code-block"> public class CollectRequest { </div> <div class="ql-code-block"> 描述该属性 </div> <div class="ql-code-block"> @ApiModelProperty("收藏的视频的id") </div> <div class="ql-code-block"> private Integer videoId; </div> <div class="ql-code-block"> @ApiModelProperty("该用户的收藏夹的id") </div> <div class="ql-code-block"> private Integer collectGroupId; </div> <div class="ql-code-block"> @ApiModelProperty("操作类型,true是收藏,false是取消收藏") </div> <div class="ql-code-block"> private Boolean type; </div> <div class="ql-code-block"> } </div> </div> <p>接口中描述信息</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @RestController </div> <div class="ql-code-block"> 给接口文档中的该接口起名,否则默认是接口类的名字CollectController </div> <div class="ql-code-block"> @Api(tags = "收藏和收藏夹的增删改查") </div> <div class="ql-code-block"> @RequestMapping("/collect") </div> <div class="ql-code-block"> public class CollectController { </div> <div class="ql-code-block"> @Resource </div> <div class="ql-code-block"> CollectService collectService; </div> <div class="ql-code-block"> @PostMapping("/collect") </div> <div class="ql-code-block"> 给该接口方法添加描述 </div> <div class="ql-code-block"> @ApiOperation("收藏视频") </div> <div class="ql-code-block"> public Result<Boolean> collect(@RequestBody List<CollectRequest> collectRequest) { </div> <div class="ql-code-block"> return collectService.collect(collectRequest); </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> } </div> </div> <p>实际效果(浏览器中输入服务ip+端口+/doc.html)</p> <p><img src="https://pic.code-nav.cn/planet_post_image/1808693273146032129/yn3fqp2v.jpeg"></p> <p>集成knife4j后页面好看很多,且封装了一些参数,操作更便捷。原生swagger是这样</p> <p><img src="https://pic.code-nav.cn/planet_post_image/1808693273146032129/p61zlqhf.jpeg"></p> <p>openfeign的使用</p> <p>引入依赖</p> <div class="ql-code-block-container"> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>org.springframework.cloud</groupId> </div> <div class="ql-code-block"> <artifactId>spring-cloud-starter-openfeign</artifactId> 在当前版本引入spring-cloud-starter-openfeign-core会报错,需注意不要引错依赖 </div> <div class="ql-code-block"> </dependency> </div> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>org.springframework.cloud</groupId> </div> <div class="ql-code-block"> <artifactId>spring-cloud-starter-loadbalancer</artifactId>当前版本openfeign已经弃用ribbon,因此需额外引入loadbalancer作为负载均衡器 </div> <div class="ql-code-block"> </dependency> </div> </div> <p>client接口编写</p> <div class="ql-code-block-container"> <div class="ql-code-block"> name值是服务注册到注册中心的名称,url值可以省去,若同时写则url值优先级更高,将不走注册中心直接调该路径的方法 </div> <div class="ql-code-block"> @FeignClient(name = "notice",url="http://localhost:30000") </div> <div class="ql-code-block"> @Component </div> <div class="ql-code-block"> public interface SendNoticeClient { </div> <div class="ql-code-block"> 需和实际远程调用接口路径、返回值、参数列表一致 </div> <div class="ql-code-block"> @PostMapping("/notice/sendDynamicNotice") </div> <div class="ql-code-block"> Boolean dynamicNotice(@RequestBody Dynamic dynamic); </div> </div> <p>使用</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @Resource </div> <div class="ql-code-block"> SendNoticeClient client; </div> <div class="ql-code-block"> client.dynamicNotice(new Dynamic()); </div> </div> <p>绝大多数情况传递的都是json数据,但特殊情况需要在服务之间传递文件流,这时需要额外的措施</p> <p>调用的方法传参中有文件</p> <p>由于MultipartFile接口往往是前端上传文件后由springboot封装好的流程绑定到后端接口上,无法实例化,因此需要自定义一个类实现MultipartFile接口</p> <div class="ql-code-block-container"> <div class="ql-code-block"> public class CustomMultipartFile implements MultipartFile { </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> private final byte[] fileContent; </div> <div class="ql-code-block"> private final String fileName; </div> <div class="ql-code-block"> private final String contentType; </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> public CustomMultipartFile(InputStream inputStream, String fileName, String contentType) throws IOException { </div> <div class="ql-code-block"> this.fileContent = IoUtil.readBytes(inputStream); </div> <div class="ql-code-block"> this.fileName = fileName; </div> <div class="ql-code-block"> this.contentType = contentType; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> public CustomMultipartFile(byte[] fileContent, String fileName, String contentType) throws IOException { </div> <div class="ql-code-block"> this.fileContent = fileContent; </div> <div class="ql-code-block"> this.fileName = fileName; </div> <div class="ql-code-block"> this.contentType = contentType; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public String getName() { </div> <div class="ql-code-block"> return StringUtils.cleanPath(fileName); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public String getOriginalFilename() { </div> <div class="ql-code-block"> return StringUtils.cleanPath(fileName); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public String getContentType() { </div> <div class="ql-code-block"> return contentType; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public boolean isEmpty() { </div> <div class="ql-code-block"> return fileContent == null || fileContent.length == 0; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public long getSize() { </div> <div class="ql-code-block"> return fileContent.length; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public byte[] getBytes() { </div> <div class="ql-code-block"> return fileContent; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public InputStream getInputStream() { </div> <div class="ql-code-block"> InputStream inputStream=IoUtil.toStream(getBytes()); </div> <div class="ql-code-block"> return inputStream; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"><br> </div> <div class="ql-code-block"> @Override </div> <div class="ql-code-block"> public void transferTo(java.io.File dest) throws IllegalStateException { </div> <div class="ql-code-block"> throw new UnsupportedOperationException("transfer to file not supported"); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <p>client和实际远程调用接口都需要额外添加注解中cosumes的值</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @PostMapping(value = "/videoEncode/uploadVideo",consumes = MediaType.MULTIPART_FORM_DATA_VALUE) </div> <div class="ql-code-block"> void uploadVideo(@RequestPart("multipartFile") MultipartFile multipartFile); </div> </div> <p>这样实际调用时可以传入实例化对象</p> <div class="ql-code-block-container"> <div class="ql-code-block"> CustomMultipartFile customMultipartFile = new CustomMultipartFile(inputStream, uploadVideo.getUrl().substring(uploadVideo.getUrl().lastIndexOf("/") + 1) </div> <div class="ql-code-block"> , contentType); </div> <div class="ql-code-block"> videoClient.uploadVideo(customMultipartFile); </div> </div> <p>响应是文件</p> <p>需要用到InputstreamResource类进行装载</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @PostMapping("/getVideoInputStream") </div> <div class="ql-code-block"> public ResponseEntity<Resource> getVideoInputStream(@RequestBody UploadVideo uploadVideo){ </div> <div class="ql-code-block"> String url=uploadVideo.getUrl(); </div> <div class="ql-code-block"> InputStream inputStream = minioService.getObject(url.substring(url.lastIndexOf("/")+1)); </div> <div class="ql-code-block"> 将文件流封装进InputStreamResource对象中 </div> <div class="ql-code-block"> InputStreamResource resource = new InputStreamResource(inputStream); </div> <div class="ql-code-block"> 将InputStreamResource对象封装进远程调用得到的响应的响应体中 </div> <div class="ql-code-block"> return ResponseEntity.ok() </div> <div class="ql-code-block"> .header(HttpHeaders.CONTENT_DISPOSITION, HEADERS_VALUES) </div> <div class="ql-code-block"> .body(resource); </div> <div class="ql-code-block"> } </div> </div> <p>client接口写法同传统json传递方法的写法,只是返回响应类型比较特殊</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @PostMapping("/videoEncode/getVideoInputStream") </div> <div class="ql-code-block"> ResponseEntity<Resource> getVideo(@RequestBody UploadVideo uploadVideo); </div> </div> <p>最后,宣传一下自己的<span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">仿b站前后端分离微服务项目,依赖版本号也在该项目的父pom.xml中</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">实现了以下功能:</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">视频的上传、查看与上传时获取封面</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">视频的点赞、评论、可同时新增和删除多个收藏记录的收藏、多功能的弹幕</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">用户的个人信息查看编辑、用户之间的关注</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">用户的个人主页权限修改、查看、由个人主页权限动态决定的用户个人主页内容的获取</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">手机号、邮箱、图形验证码的多种方式登录</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">支持临时会话的服务器为代理的一对一实时私聊</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">基于讯飞星火的文生文、文生图、(全网首发)智能PPT</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">关注up动态视频、评论、点赞、私聊消息的生成与推送</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">基于es实现的视频和用户的聚合搜索、推荐视频</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">网关的路由和统一鉴权与授权</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">基于双token的七天内无感刷新token</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">防csrf、xss、抓包、恶意上传脚本攻击</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">统一处理异常和泛型封装响应体、自定义修改响应序列化值</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">简易的仿redis缓存读取与数据过期剔除实现</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">xxl-job+ redis+ rocketmq+ es+ 布隆过滤器的自定义es与mysql数据同步</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">slueth+zipkin的多服务间请求链路追踪</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">集中多服务日志到一个文件目录下与按需添加特定内容入日志</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">多服务的详细接口文档</span></p> <p><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">项目地址</span><a href="https://labilibili.com/" target="_blank" style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(86, 120, 149);">LABiliBili</a><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">,github地址</span><a href="https://github.com/aigcbilibili/aigcbilibili" target="_blank" style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(86, 120, 149);">GitHub - aigcbilibili/aigcbilibili: 仿bilibili前后端实现</a><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">,演示地址</span><a href="https://labilibili.com/video/%E6%BC%94%E7%A4%BA.mp4" target="_blank" style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(86, 120, 149);">https://labilibili.com/video/</a><a href="https://labilibili.com/video/%E6%BC%94%E7%A4%BA.mp4" target="_blank" style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">演示.mp4</a><span style="background-color: rgb(255, 255, 255); font-size: 15px; color: rgb(47, 48, 52);">,如果大家觉得有帮助的话可以去github点个小星星♪(・ω・)ノ</span></p> <p><br></p> <p><br></p> </div> </body> </html>
WebSocket 在 JS 中的使用以及在 SpringBoot 中整合 WebSocket
<html> <head></head> <body> <div class="content ql-editor"> <h1>前端 WebSocket 的一些使用</h1> <p>WebSocket 是一种网络通信协议,用于实现双向通信。在前端中,你可以使用 JavaScript 中的 <code>WebSocket</code> 对象来创建 WebSocket 连接,发送和接收数据。</p> <h2>连接的建立</h2> <p>通过创建一个 WebSocket 对象建立一个 WebSocket 连接</p> <p>例如:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> const ws = new WebSocket('ws://localhost:8080/channel/echo'); </div> </div> <p>传给对象的参数是通过 WebSocket 协议通讯的网络地址。</p> <h2>接收消息</h2> <p>接收消息这里指的是接收服务端的消息。</p> <p>这里有两种方法。</p> <ol> <li data-list="ordered"><span class="ql-ui"></span><strong>使用 <code>addEventListener</code></strong>: 你可以使用 <code>addEventListener</code> 来监听 <code>message</code> 事件,这是最常见的方式,也是推荐的做法。 <div class="ql-code-block-container"> <div class="ql-code-block"> ws.addEventListener('message', (event) => { const receivedMessage = event.data; console.log('Received message from server:', receivedMessage); // 在这里处理接收到的消息 }); </div> </div></li> <li data-list="ordered"><span class="ql-ui"></span><strong>使用 <code>onmessage</code> 属性</strong>: 除了使用 <code>addEventListener</code>,你还可以直接设置 <code>onmessage</code> 属性来指定消息处理函数。这与之前的示例相似,但更简洁: <div class="ql-code-block-container"> <div class="ql-code-block"> ws.onmessage = function (event) { const receivedMessage = event.data; console.log('Received message from server:', receivedMessage); // 在这里处理接收到的消息 }; </div> </div></li> </ol> <h2>发送消息</h2> <p><strong>发送消息到服务器</strong>: 使用 <code>send()</code> 方法将消息发送到服务器:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> ws.send('Hello, server!'); // 发送消息给服务器 </div> </div> <h2>关闭连接</h2> <p><strong>关闭 WebSocket 连接</strong>: 要关闭 WebSocket 连接,你可以简单地使用 <code>WebSocket.close()</code> 方法,例如:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> ws.close(); </div> </div> <p>如果 WebSocket 连接的 <code>readyState</code> 已经处于 <code>CLOSE</code> 状态,那么该方法不会执行任何操作</p> <p>检查 WebSocket 是否打开: 你可以通过检查 <code>WebSocket</code> 的 <code>readyState</code> 属性来判断 WebSocket 是否打开。如果 <code>readyState</code> 的值为 <code>WebSocket.OPEN</code>,则表示连接已打开:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> if (ws.readyState === WebSocket.OPEN) { // WebSocket 连接已打开 } </div> </div> <p>这样你就可以在代码中判断 WebSocket 是否处于打开状态了</p> <h2>处理</h2> <p><strong>处理连接状态</strong>: 你可以监听其他事件,例如 <code>open</code>、<code>close</code> 和 <code>error</code>,以处理连接的不同状态:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> ws.addEventListener('open', (event) => { console.log('WebSocket 已连接'); }); ws.addEventListener('close', (event) => { console.log('WebSocket 连接已关闭'); }); ws.addEventListener('error', (event) => { console.error('WebSocket 连接出现异常:', event.error); }); </div> </div> <p>同样可以使用onclose 、 onerror 、 onopen 属性定义时间监听函数。</p> <h1>在 Spring Boot 中整合、使用 WebSocket</h1> <p>WebSocket 是一种基于 TCP 协议的全双工通信协议,它允许客户端和服务器之间建立持久的、双向的通信连接。相比传统的 HTTP 请求 - 响应模式,WebSocket 提供了实时、低延迟的数据传输能力。通过 WebSocket,客户端和服务器可以在任意时间点互相发送消息,实现实时更新和即时通信的功能。WebSocket 协议经过了多个浏览器和服务器的支持,成为了现代 Web 应用中常用的通信协议之一。它广泛应用于聊天应用、实时数据更新、多人游戏等场景,为 Web 应用提供了更好的用户体验和更高效的数据传输方式。</p> <p>本文将会指导你如何在 Spring Boot 中整合、使用 WebSocket,以及如何在 <code>@ServerEndpoint</code> 类中注入其他 Bean 依赖 。</p> <p>在 Spring Boot 中使用 WebSocket 有 2 种方式。第 1 种是使用由 Jakarta EE 规范提供的 Api,也就是 <code>jakarta.websocket</code> 包下的接口。第 2 种是使用 spring 提供的支持,也就是 <a href="https://github.com/spring-projects/spring-framework/tree/main/spring-websocket" target="_blank"><code>spring-websocket</code></a> 模块。前者是一种独立于框架的技术规范,而后者是 Spring 生态系统的一部分,可以与其他 Spring 模块(如 Spring MVC、Spring Security)无缝集成,共享其配置和功能。</p> <p>2 种方式各有优劣,你可以按需选择。本文将使用第 1 种方式,也就是使用 <code>jakarta.websocket</code> 来开发 WebSocket 应用。</p> <p>软件版本:</p> <ol> <li data-list="bullet"><span class="ql-ui"></span>Spring Boot:<code>3.1.3</code></li> </ol> <h2>在 Spring Boot 中整合 WebSocket</h2> <h3>添加依赖</h3> <p>在 <code>pom.xml</code> 中添加 <code>spring-boot-starter-websocket</code> 依赖。</p> <div class="ql-code-block-container"> <div class="ql-code-block"> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency> </div> </div> <h3>开发 ServerEndpoint 端点</h3> <p>服务端 WebSocket 端点的开发也有 2 种方式。第 1 种是实现规范所提供的各种接口,通过接口定义的回调方法来处理新的连接、客户端消息、连接断开等等事件。另一种方式是使用注解,类似于 Spring 中的 Controller,通过在方法上使用不同的注解来监听不同的 WebSocket 事件,灵活性比较高,推荐使用。</p> <p>我们打算创建一个 <code>echo</code> 端点,该端点会处理客户端的连接、断开、消息事件。在收到消息后,我们会在消息前面加上服务器时间戳和 <code>Hello</code> 前缀,原样写回给客户端。如果客户端发送的消息为 <code>bye</code>,则服务器会主动断开与客户端的连接。</p> <div class="ql-code-block-container"> <div class="ql-code-block"> package cn.springdoc.demo.channel; import java.io.IOException; import java.time.Instant; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import jakarta.websocket.CloseReason; import jakarta.websocket.EndpointConfig; import jakarta.websocket.OnClose; import jakarta.websocket.OnError; import jakarta.websocket.OnMessage; import jakarta.websocket.OnOpen; import jakarta.websocket.Session; import jakarta.websocket.server.ServerEndpoint; // 使用 @ServerEndpoint 注解表示此类是一个 WebSocket 端点 // 通过 value 注解,指定 websocket 的路径 @ServerEndpoint(value = "/channel/echo") public class EchoChannel { private static final Logger LOGGER = LoggerFactory.getLogger(EchoChannel.class); private Session session; // 收到消息 @OnMessage public void onMessage(String message) throws IOException{ LOGGER.info("[websocket] 收到消息:id={},message={}", this.session.getId(), message); if (message.equalsIgnoreCase("bye")) { // 由服务器主动关闭连接。状态码为 NORMAL_CLOSURE(正常关闭)。 this.session.close(new CloseReason(CloseReason.CloseCodes.NORMAL_CLOSURE, "Bye"));; return; } this.session.getAsyncRemote().sendText("["+ Instant.now().toEpochMilli() +"] Hello " + message); } // 连接打开 @OnOpen public void onOpen(Session session, EndpointConfig endpointConfig){ // 保存 session 到对象 this.session = session; LOGGER.info("[websocket] 新的连接:id={}", this.session.getId()); } // 连接关闭 @OnClose public void onClose(CloseReason closeReason){ LOGGER.info("[websocket] 连接断开:id={},reason={}", this.session.getId(),closeReason); } // 连接异常 @OnError public void onError(Throwable throwable) throws IOException { LOGGER.info("[websocket] 连接异常:id={},throwable={}", this.session.getId(), throwable.getMessage()); // 关闭连接。状态码为 UNEXPECTED_CONDITION(意料之外的异常) this.session.close(new CloseReason(CloseReason.CloseCodes.UNEXPECTED_CONDITION, throwable.getMessage())); } } </div> </div> <p>首先,使用 <code>@ServerEndpoint</code> 注解表示此类是一个 WebSocket 端点,<code>value</code> 属性是必须的,用于设置路由。它还有其他的一些可选属性可以用于自定义子协议、消息编码器、消息解码器、握手处理器等等,篇幅原因这里不展开。</p> <h4>@OnMessage</h4> <p><code>@OnMessage</code> 注解用于监听客户端消息事件,它只有一个属性 <code>long maxMessageSize() default -1;</code> 用于限制客户端消息的大小,如果小于等于 0 则表示不限制。当客户端消息体积超过这个阈值,那么服务器就会主动断开连接,状态码为:<code>1009</code>。方法的参数可以是基本的 <code>String</code> / <code>byte[]</code> 或者是 <code>Reader</code> / <code>InputStream</code>,分别表示 WebSocket 中的文本和二进制消息。也可以是自定义的 Java 对象,但是需要在 <code>@ServerEndpoint</code> 中配置对象的解码器(<code>jakarta.websocket.Decoder</code>)。对于内容较长的消息,支持分批发送,可以在消息参数后面定义一个布尔类型的 <code>boolean last</code>参数,如果该值为 <code>true</code> 则表示此消息是批次消息中的最后一条。</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @OnMessage public void onMessage(String message, boolean last) throws IOException{ if (last) { // 这是批量消息的最后一条 } } </div> </div> <h4>@OnOpen</h4> <p><code>@OnOpen</code> 方法用于监听客户端的连接事件,它没有任何属性。可以作为方法参数的对象有很多,<code>Session</code> 对象是必须的,表示当前连接对象,我们可以通过此对象来执行发送消息、断开连接等操作。WebSocket 的连接 URL,类似于 Http 的 URL,也可以传递查询参数、path 参数。通常用于传递认证、鉴权用的 Token 或其他信息。</p> <p>要获取查询参数,我们可以通过 <code>Session</code> 的 <code>getRequestParameterMap();</code> 获取。</p> <div class="ql-code-block-container"> <div class="ql-code-block"> Map<String, List<String>> query = session.getRequestParameterMap(); </div> </div> <p>要获取 path 参数,首先要在 <code>@ServerEndpoint</code> 中定义 path 参数,类似于 Spring Mvc 的 path 参数定义。例如: <code>@ServerEndpoint(value = "/channel/echo/{id}")</code>。那么我们可以在 <code>@OnOpen</code> 方法中使用 <code>@PathParam</code> 注解接收,如下:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @ServerEndpoint(value = "/channel/echo/{id}") ... @OnOpen public void onOpen(Session session, @PathParam("id") Long id, EndpointConfig endpointConfig){ .... } </div> </div> <p>示例中的最后一个参数 <code>EndpointConfig</code> ,它是可选,用于获取全局的一些配置。在本文中未用到。</p> <h4>@OnClose</h4> <p><code>@OnClose</code> 用于处理连接断开事件,参数中可以指定一个 <code>CloseReason</code> 对象,它封装了断开连接的状态码、原因信息。</p> <h4>@OnError</h4> <p><code>@OnError</code> 用于处理异常事件,<strong>该方法必须要有一个 <code>Throwable</code> 类型的参数</strong>,表示发生的异常。否则应用会启用失败:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> Caused by: jakarta.websocket.DeploymentException: No Throwable parameter was present on the method [onError] of class [cn.springdoc.demo.channel.EchoChannel] that was annotated with OnError at org.apache.tomcat.websocket.pojo.PojoMethodMapping.getPathParams(PojoMethodMapping.java:311) ~[tomcat-embed-websocket-10.1.12.jar:10.1.12] at org.apache.tomcat.websocket.pojo.PojoMethodMapping.<init>(PojoMethodMapping.java:194) ~[tomcat-embed-websocket-10.1.12.jar:10.1.12] at org.apache.tomcat.websocket.server.WsServerContainer.addEndpoint(WsServerContainer.java:130) ~[tomcat-embed-websocket-10.1.12.jar:10.1.12] at org.apache.tomcat.websocket.server.WsServerContainer.addEndpoint(WsServerContainer.java:240) ~[tomcat-embed-websocket-10.1.12.jar:10.1.12] at org.apache.tomcat.websocket.server.WsServerContainer.addEndpoint(WsServerContainer.java:198) ~[tomcat-embed-websocket-10.1.12.jar:10.1.12] at org.springframework.web.socket.server.standard.ServerEndpointExporter.registerEndpoint(ServerEndpointExporter.java:156) ~[spring-websocket-6.0.11.jar:6.0.11] ... 12 common frames omitted </div> </div> <p>所有事件方法,都支持使用 <code>Session</code> 作为参数,表示当前连接参数。但是为了更加方便,我们在 <code>@OnOpen</code> 事件中直接把 <code>Session</code> 存储到了当前对象中,可以在任意方法中使用 <code>this</code> 访问。服务器会为每个连接创建一个端点对象,所以这是线程安全的。</p> <p>上面还提到了一个 “连接关闭状态码”,WebSocket 协议定义了一系列状态码来表示连接断开的原因,这些状态码定义在了 <code>CloseReason.CloseCodes</code> 枚举中。</p> <h3>配置 ServerEndpointExporter</h3> <p>定义好端点后,需要在配置类中通过定义 <code>ServerEndpointExporter</code> Bean 进行注册。</p> <div class="ql-code-block-container"> <div class="ql-code-block"> package cn.springdoc.demo.config; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.web.socket.server.standard.ServerEndpointExporter; import cn.springdoc.demo.channel.EchoChannel; @Configuration public class WebSocketConfiguration { @Bean public ServerEndpointExporter serverEndpointExporter (){ ServerEndpointExporter exporter = new ServerEndpointExporter(); // 手动注册 WebSocket 端点 exporter.setAnnotatedEndpointClasses(EchoChannel.class); return exporter; } } </div> </div> <p>你也可以在 WebSocket 端点上添加 <code>@Component</code> 注解,使用 Spring 自动扫描,这样的话不需要手动调用 <code>setAnnotatedEndpointClasses</code> 方法进行注册。</p> <h2>测试</h2> <p>在项目的 <code>src/main/resources</code> 目录下创建一个 <code>public</code> 文件夹,再在此文件夹中新建一个 <code>index.html</code> 文件,作为 WebSocket 客户端。内容如下:</p> <blockquote> <p>Spring Boot 默认会把 <code>public</code> 目录下的 <code>index.html</code> 作为应用主页。</p> </blockquote> <div class="ql-code-block-container"> <div class="ql-code-block"> <!DOCTYPE html> <html> <head> <meta charset="UTF-8"> <title>WebSocket</title> </head> <body> <script type="text/javascript"> let websocket = new WebSocket("ws://localhost:8080/channel/echo"); // 连接断开 websocket.onclose = e => { console.log(`连接关闭: code=${e.code}, reason=${e.reason}`) } // 收到消息 websocket.onmessage = e => { console.log(`收到消息:${e.data}`); } // 异常 websocket.onerror = e => { console.log("连接异常") console.error(e) } // 连接打开 websocket.onopen = e => { console.log("连接打开"); // 创建连接后,往服务器连续写入3条消息 websocket.send("sprigdoc.cn"); websocket.send("sprigdoc.cn"); websocket.send("sprigdoc.cn"); // 最后发送 bye,由服务器断开连接 websocket.send("bye"); // 也可以由客户端主动断开 // websocket.close(); } </script> </body> </html> </div> </div> <p>内容很简单,网页加载后运行 Javascript 代码。立即创建与 <code>ws://localhost:8080/channel/echo</code> 的 WebSocket 连接对象,通过注册对象的各种监听方法来处理事件。</p> <p>在连接就绪后,也就是在 <code>onopen</code> 方法中往服务器端点发送了 3 条消息。按照逻辑,服务端也会回复 3 条消息,这会触发 <code>onmessage</code> 事件,把消息内容输出到控制台。最后,发送 <code>bye</code>,服务器收到消息后会主动断开连接,这就会触发 <code>onclose</code> 事件,把 “连接关闭状态码” 和原因输出到控制台。</p> <blockquote> <p>其实你可以直接把这段 Javascript 代码复制到任意支持 WebSocket 的浏览器的控制台执行,WebSocket 没有跨域的说法!</p> </blockquote> <p>启动应用,打开浏览器(先打开控制台),然后访问 <code>http://localhost:8080/</code>,查看控制台输出的日志:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> 连接打开 收到消息:[1694505275009] Hello sprigdoc.cn 收到消息:[1694505275012] Hello sprigdoc.cn 收到消息:[1694505275014] Hello sprigdoc.cn 连接关闭: code=1000, reason=Bye </div> </div> <p>再看看服务端控制台日志:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> cn.springdoc.demo.channel.EchoChannel : [websocket] 新的连接:id=0 cn.springdoc.demo.channel.EchoChannel : [websocket] 收到消息:id=0,message=sprigdoc.cn cn.springdoc.demo.channel.EchoChannel : [websocket] 收到消息:id=0,message=sprigdoc.cn cn.springdoc.demo.channel.EchoChannel : [websocket] 收到消息:id=0,message=sprigdoc.cn cn.springdoc.demo.channel.EchoChannel : [websocket] 收到消息:id=0,message=bye cn.springdoc.demo.channel.EchoChannel : [websocket] 连接断开:id=0,reason=CloseReason: code [1000], reason [Bye] </div> </div> <p>没有任何问题,一切按照我们预定义的逻辑在运行。客户端发送 3 条消息,服务器响应 3 条消息,最后断开连接。客户端、服务器相应的事件方法都成功执行。</p> <p>服务端日志中的 sessionId(<code>id=0</code>),是通过 <code>Session</code> 的 <code>String getId();</code> 方法获取的。服务器会为每个连接分配一个不同的 id 值,不同服务器生成的 id 类型不一样。 Tomcat 使用从 0 开始的自增值(本例),Undertow 使用的是类似于 UUID 的 32 位长度的字符串。</p> <h2>在端点中注入 Bean</h2> <p>往往我们需要在端点中使用其他 Spring 管理的 Bean 来完成业务,例如认证、鉴权、保存消息。。。等等。</p> <p>假如我们有一个 <code>UserService</code> 服务类,内容如下:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> package cn.springdoc.demo.service; import org.springframework.stereotype.Service; @Service public class UserService { public void foo() {} // .... } </div> </div> <p>我们现在要在端点中注入使用它,很多人会直接在端点类上使用 <code>@Component</code> 注解,然后注入:</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @ServerEndpoint(value = "/channel/echo") @Component // 注册为 Spring 组件 public class EchoChannel { @Autowired // 注入需要的 Bean private UserService userService; // ... @OnOpen public void onOpen(Session session, EndpointConfig endpointConfig){ this.session = session; // 在业务中使用 this.userService.foo(); } } </div> </div> <p>服务可以正常启动,看似一切都没问题!可是当你在事件方法中使用这 Bean 的时候就会导致 <code>NullPointerException</code> 异常。</p> <div class="ql-code-block-container"> <div class="ql-code-block"> java.lang.NullPointerException: Cannot invoke "cn.springdoc.demo.service.UserService.foo()" because "this.userService" is null at cn.springdoc.demo.channel.EchoChannel.onOpen(EchoChannel.java:54) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:568) at org.apache.tomcat.websocket.pojo.PojoEndpointBase.doOnOpen(PojoEndpointBase.java:67) at org.apache.tomcat.websocket.pojo.PojoEndpointServer.onOpen(PojoEndpointServer.java:46) at org.apache.tomcat.websocket.server.WsHttpUpgradeHandler.init(WsHttpUpgradeHandler.java:131) at org.apache.coyote.AbstractProtocol$ConnectionHandler.process(AbstractProtocol.java:936) at org.apache.tomcat.util.net.NioEndpoint$SocketProcessor.doRun(NioEndpoint.java:1740) at org.apache.tomcat.util.net.SocketProcessorBase.run(SocketProcessorBase.java:52) at org.apache.tomcat.util.threads.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1191) at org.apache.tomcat.util.threads.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:659) at org.apache.tomcat.util.threads.TaskThread$WrappingRunnable.run(TaskThread.java:61) at java.base/java.lang.Thread.run(Thread.java:833) </div> </div> <p><strong>原因:运行时的 WebSocket 连接对象,也就是端点实例,是由服务器创建,而不是 Spring,所以不能使用自动装配</strong>。上文也提到过 “服务器会为每个连接创建一个端点实例对象”。</p> <p>知道了原因后,解决办法也很简单,我们可以使用 Spring 的 <code>ApplicationContextAware</code> 接口,在应用启动时获取到 <code>ApplicationContext</code> 并且保存在全局静态变量中。</p> <p>服务器每次创建连接的时候,我们就在 <code>@OnOpen</code> 事件方法中从 <code>ApplicationContext</code> 获取到需要 Bean 来初始化端点对象。</p> <div class="ql-code-block-container"> <div class="ql-code-block"> @ServerEndpoint(value = "/channel/echo") @Component // 由 spring 扫描管理 public class EchoChannel implements ApplicationContextAware { // 实现 ApplicationContextAware 接口, Spring 会在运行时注入 ApplicationContext private static final Logger LOGGER = LoggerFactory.getLogger(EchoChannel.class); // 全局静态变量,保存 ApplicationContext private static ApplicationContext applicationContext; private Session session; // 声明需要的 Bean private UserService userService; // 保存 Spring 注入的 ApplicationContext 到静态变量 @Override public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { EchoChannel.applicationContext = applicationContext; } @OnOpen public void onOpen(Session session, EndpointConfig endpointConfig){ // 保存 session 到对象 this.session = session; // 连接创建的时候,从 ApplicationContext 获取到 Bean 进行初始化 this.userService = EchoChannel.applicationContext.getBean(UserService.class); // 在业务中使用 this.userService.foo(); LOGGER.info("[websocket] 新的连接:id={}", this.session.getId()); } // .... } </div> </div> <p><code>onOpen</code> 方法在整个连接的生命周期中,只会执行一次,所以这种方式不会带来通信时的性能损耗。</p> #Java# #WebSocket# #开发交流# </div> </body> </html>
推送数据?也许你不需要 WebSocket
<html> <head></head> <body> <div class="content ql-editor"> <p>提到推送数据,大家可能会首先想到 WebSocket。</p> <p></p> <p>确实,WebSocket 能双向通信,自然也能做服务器到浏览器的消息推送。</p> <p></p> <p>但如果只是单向推送消息的话,HTTP 就有这种功能,它就是 Server Send Event。</p> <p></p> <p>WebSocket 的通信过程是这样的:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/dwj7i5q7.jpeg" alt=""></p> <p></p> <p>首先通过 http 切换协议,服务端返回 101 的状态码后,就代表协议切换成功。</p> <p></p> <p>之后就是 WebSocket 格式数据的通信了,一方可以随时向另一方推送消息。</p> <p></p> <p>而 HTTP 的 Server Send Event 是这样的:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/ynqhj625.jpeg" alt=""></p> <p></p> <p>服务端返回的 Content-Type 是 text/event-stream,这是一个流,可以多次返回内容。</p> <p></p> <p>Sever Send Event 就是通过这种消息来随时推送数据。</p> <p></p> <p>可能你是第一次听说 SSE,但你肯定用过基于它的应用。</p> <p></p> <p>比如你用的 CICD 平台,它的日志是实时打印的。</p> <p></p> <p>那它是如何实时传输构建日志的呢?</p> <p></p> <p>明显需要一段一段的传输,这种一般就是用 SSE 来推送数据。</p> <p></p> <p>再比如说 ChatGPT,它回答一个问题不是一次性给你全部的,而是一部分一部分的加载回答。</p> <p></p> <p>这也是基于 SSE。</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/o7rtk556.jpeg" alt=""></p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/fpdgsvkj.jpeg" alt=""></p> <p></p> <p>知道了什么是 SSE 以及它的应用,我们来自己实现一下吧:</p> <p></p> <p>创建 nest 项目:</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"> npx nest <span class="ql-token hljs-keyword">new</span> <span class="ql-token hljs-title">sse</span>-test </div> </div> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/c9psb2z6.jpeg" alt=""></p> <p></p> <p>把它跑起来:</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"> npm run start:dev </div> </div> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/dxcdgh02.jpeg" alt=""></p> <p></p> <p>访问 http://localhost:3000 可以看到 hello world,代表服务器跑成功了:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/ul3ewgc5.jpeg" alt=""></p> <p></p> <p>然后在 AppController 添加一个 stream 接口:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/uqm8an7d.jpeg" alt=""></p> <p></p> <p>这里不是通过 @Get、@Post 等装饰器标识,而是通过 @Sse 标识这是一个 event stream 类型的接口。</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"><span class="ql-token hljs-meta">@Sse('stream')</span> stream() { <span class="ql-token hljs-keyword">return</span> <span class="ql-token hljs-keyword">new</span> <span class="ql-token hljs-title">Observable</span>((observer) => { observer.next({ data: { msg: <span class="ql-token hljs-string">'aaa'</span>} }); setTimeout(() => { observer.next({ data: { msg: <span class="ql-token hljs-string">'bbb'</span>} }); }, <span class="ql-token hljs-number">2000</span>); setTimeout(() => { observer.next({ data: { msg: <span class="ql-token hljs-string">'ccc'</span>} }); }, <span class="ql-token hljs-number">5000</span>); }); } </div> </div> <p></p> <p>返回的是一个 Observable 对象,然后内部用 observer.next 返回消息。</p> <p></p> <p>可以返回任意的 json 数据。</p> <p></p> <p>我们先返回了一个 aaa、过了 2s 返回了 bbb,过了 5s 返回了 ccc。</p> <p></p> <p>然后写个前端页面:</p> <p></p> <p>创建一个 react 项目:</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"> npx create-react-app --template=typescript sse-test-frontend </div> </div> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/g5pje8u5.jpeg" alt=""></p> <p></p> <p>在 App.tsx 里写如下代码:</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"><span class="ql-token hljs-keyword">import</span> { useEffect } from <span class="ql-token hljs-string">'react'</span>; function <span class="ql-token hljs-title">App()</span> { useEffect(() => { <span class="ql-token hljs-type">const</span> <span class="ql-token hljs-variable">eventSource</span> <span class="ql-token hljs-operator">=</span> <span class="ql-token hljs-keyword">new</span> <span class="ql-token hljs-title">EventSource</span>(<span class="ql-token hljs-string">'http://localhost:3000/stream'</span>); eventSource.onmessage = ({ data }) => { console.log(<span class="ql-token hljs-string">'New message'</span>, JSON.parse(data)); }; }, []); <span class="ql-token hljs-keyword">return</span> ( <div>hello</div> ); } export <span class="ql-token hljs-keyword">default</span> App; </div> </div> <p></p> <p>这个 EventSource 是浏览器原生 api,就是用来获取 sse 接口的响应的,它会把每次消息传入 onmessage 的回调函数。</p> <p></p> <p>我们在 nest 服务开启跨域支持:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/o9la9eec.jpeg" alt=""></p> <p></p> <p>然后把 react 项目 index.tsx 里这几行代码删掉,它会导致额外的渲染:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/p7akkso9.jpeg" alt=""></p> <p></p> <p>执行 npm run start</p> <p></p> <p>因为 3000 端口被占用了,它会跑在 3001:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/ycr89wgx.jpeg" alt=""></p> <p></p> <p>浏览器访问下:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/y34hlbaa.jpeg" alt=""></p> <p></p> <p>看到一段段的响应了没?</p> <p></p> <p>这就是 Server Send Event。</p> <p></p> <p>在 devtools 里可以看到,响应的 Content-Type 是 text/event-stream:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/p6dcwdfa.jpeg" alt=""></p> <p></p> <p>然后在 EventStream 里可以看到每一次收到的消息:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/krccmadc.jpeg" alt=""></p> <p></p> <p>这样,服务端就可以随时向网页推送消息了。</p> <p></p> <p>那它兼容性怎么样呢?</p> <p></p> <p>可以在 <a href="https://developer.mozilla.org/zh-CN/docs/Web/API/EventSource#%E6%B5%8F%E8%A7%88%E5%99%A8%E5%85%BC%E5%AE%B9%E6%80%A7" target="_blank">MDN</a> 看到:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/sbcr31ge.jpeg" alt=""></p> <p></p> <p>除了 ie、edge 外,其他浏览器都没任何兼容问题。</p> <p></p> <p>基本是可以放心用的。</p> <p></p> <p>那用在哪呢?</p> <p></p> <p>一些只需要服务端推送的场景就特别适合 Server Send Event。</p> <p></p> <p>比如这个站内信:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/flvvjru5.jpeg" alt=""></p> <p></p> <p>这种推送用 WebSocket 就没必要了,可以用 SSE 来做。</p> <p></p> <p>那连接断了怎么办呢?</p> <p></p> <p>不用担心,浏览器会自动重连。</p> <p></p> <p>这点和 WebSocket 不同,WebSocket 如果断开之后是需要手动重连的,而 SSE 不用。</p> <p></p> <p>再比如说日志的实时推送。</p> <p></p> <p>我们来测试下:</p> <p></p> <p>tail -f 命令可以实时看到文件的最新内容:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/9ekmfm3t.jpeg" alt=""></p> <p></p> <p>我们通过 child_process 模块的 exec 来执行这个命令,然后监听它的 stdout 输出:</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"> const { exec } = require(<span class="ql-token hljs-string">"child_process"</span>); <span class="ql-token hljs-type">const</span> <span class="ql-token hljs-variable">childProcess</span> <span class="ql-token hljs-operator">=</span> exec(<span class="ql-token hljs-string">'tail -f ./log'</span>); childProcess.stdout.on(<span class="ql-token hljs-string">'data'</span>, (msg) => { console.log(msg); }); </div> </div> <p></p> <p>用 node 执行它:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/vsdniqnc.jpeg" alt=""></p> <p></p> <p>然后添加一个 sse 的接口:</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"><span class="ql-token hljs-meta">@Sse('stream2')</span> stream2() { <span class="ql-token hljs-type">const</span> <span class="ql-token hljs-variable">childProcess</span> <span class="ql-token hljs-operator">=</span> exec(<span class="ql-token hljs-string">'tail -f ./log'</span>); <span class="ql-token hljs-keyword">return</span> <span class="ql-token hljs-keyword">new</span> <span class="ql-token hljs-title">Observable</span>((observer) => { childProcess.stdout.on(<span class="ql-token hljs-string">'data'</span>, (msg) => { observer.next({ data: { msg: msg.toString() }}); }) }); </div> </div> <p></p> <p>监听到新的数据之后,把它返回给浏览器。</p> <p></p> <p>浏览器连接这个新接口:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/kfsriaob.jpeg" alt=""></p> <p></p> <p>测试下:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/ykfrg71w.jpeg" alt=""></p> <p></p> <p>可以看到,浏览器收到了实时的日志。</p> <p></p> <p>很多构建日志都是通过 SSE 的方式实时推送的。</p> <p></p> <p>日志之类的只是文本,那如果是二进制数据呢?</p> <p></p> <p>二进制数据在 node 里是通过 Buffer 存储的。</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"> const { readFileSync } = require(<span class="ql-token hljs-string">"fs"</span>); <span class="ql-token hljs-type">const</span> <span class="ql-token hljs-variable">buffer</span> <span class="ql-token hljs-operator">=</span> readFileSync(<span class="ql-token hljs-string">'./package.json'</span>); console.log(buffer); </div> </div> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/vfld2i3n.jpeg" alt=""></p> <p></p> <p>而 Buffer 有个 toJSON 方法:</p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/oot04pew.jpeg" alt=""></p> <p></p> <p>这样不就可以通过 sse 的接口返回了么?</p> <p></p> <p>试一下:</p> <p></p> <div class="ql-code-block-container"> <div class="ql-code-block"><span class="ql-token hljs-meta">@Sse('stream3')</span> stream3() { <span class="ql-token hljs-keyword">return</span> <span class="ql-token hljs-keyword">new</span> <span class="ql-token hljs-title">Observable</span>((observer) => { <span class="ql-token hljs-type">const</span> <span class="ql-token hljs-variable">json</span> <span class="ql-token hljs-operator">=</span> readFileSync(<span class="ql-token hljs-string">'./package.json'</span>).toJSON(); observer.next({ data: { msg: json }}); }); } </div> </div> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/x2cdonka.jpeg" alt=""></p> <p></p> <p><img src="https://pic.code-nav.cn/planet_post_image/1797461877294809089/i54iayub.jpeg" alt=""></p> <p></p> <p>确实可以。</p> <p></p> <p>也就是说,基于 sse,除了可以推送文本外,还可以推送任意二进制数据。</p> <p></p> <h2>总结</h2> <p></p> <p>服务端实时推送数据,除了用 WebSocket 外,还可以用 HTTP 的 Server Send Event。</p> <p></p> <p>只要 http 返回 Content-Type 为 text/event-stream 的 header,就可以通过 stream 的方式多次返回消息了。</p> <p></p> <p>它传输的是 json 格式的内容,可以用来传输文本或者二进制内容。</p> <p></p> <p>我们通过 Nest 实现了 sse 的接口,用 @Sse 装饰器标识方法,然后返回 Observe 对象就可以了。内部可以通过 observer.next 随时返回数据。</p> <p></p> <p>前端使用 EventSource 的 onmessage 来接收消息。</p> <p></p> <p>这个 api 的兼容性很好,除了 ie 外可以放心的用。</p> <p></p> <p>它的应用场景有很多,比如站内信、构建日志实时展示、chatgpt 的消息返回等。</p> <p></p> <p>再遇到需要消息推送的场景,不要直接 WebSocket 了,也许 Server Send Event 更合适呢?</p> <p>#前端# #后端#</p> </div> </body> </html>
WebSocket的使用Demo
<html> <head></head> <body> <div class="content ql-editor"> <p>新手博客人上路,如果星球里的md格式不好的话,可以来我的博客看一看</p> <p><a href="https://adagio7_5.gitee.io/adagio_blog/2023/08/29/WebSocket%E7%9A%84%E4%BD%BF%E7%94%A8Demo/" target="_blank">WebSocket的使用Demo | Adagio Blog (</a><a href="http://gitee.io" target="_blank">gitee.io</a><a href="https://adagio7_5.gitee.io/adagio_blog/2023/08/29/WebSocket%E7%9A%84%E4%BD%BF%E7%94%A8Demo/" target="_blank">)</a></p> <p><br></p> <h1><strong style="color: rgb(52, 73, 94);">WebSocket的使用Demo</strong></h1> <h2><strong style="color: rgb(52, 73, 94);">起因</strong></h2> <p><br></p> <p>周日一个朋友问我,当用户扫描二维码进行核验之后,如何自动刷新前端页面呢,我说前端可以发请求等待回调啊,后来想想好像不对,扫码这个操作,是用户发起的,并不是前端发起的,那前端肯定是不能监听到核验这个操作的结果的,在这个过程里,似乎前端才是服务端,后端是客户端,后端确定核验之后向前端发请求,这似乎有悖于我之前的理解</p> <p>后来上网查了一下,HTML5推出了WebSocket技术,最主要的功能是可以让浏览器有着双向通信的能力,这不就很符合“前端是服务端”的需求嘛,于是我就去学习了一下WebSocket的用法,并仿写了一个小Demo</p> <p><br></p> <h2><strong style="color: rgb(52, 73, 94);">实现</strong></h2> <h3><strong style="color: rgb(52, 73, 94);">具体逻辑</strong></h3> <p><br></p> <p>前端与后端建立长连接,后端如果发生了订单状态变化,就向前端发起请求,前端进行页面刷新</p> <p>ps:并没有实际订单和其他更严谨的逻辑,主打一个模拟</p> <p><br></p> <h3><strong style="color: rgb(52, 73, 94);">后端的实现</strong></h3> <h4><strong style="color: rgb(52, 73, 94);">依赖</strong></h4> <p><br></p> <p>用的就是常见的SpringBoot框架和WebSocket的Starter</p> <div class="ql-code-block-container"> <div class="ql-code-block"> <dependencies> </div> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>org.springframework.boot</groupId> </div> <div class="ql-code-block"> <artifactId>spring-boot-starter</artifactId> </div> <div class="ql-code-block"> <version>2.4.2</version> </div> <div class="ql-code-block"> </dependency> </div> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>org.projectlombok</groupId> </div> <div class="ql-code-block"> <artifactId>lombok</artifactId> </div> <div class="ql-code-block"> <version>1.18.26</version> </div> <div class="ql-code-block"> </dependency> </div> <div class="ql-code-block"> <!--WebSocket依赖--> </div> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>org.springframework.boot</groupId> </div> <div class="ql-code-block"> <artifactId>spring-boot-starter-websocket</artifactId> </div> <div class="ql-code-block"> <version>2.4.2</version> </div> <div class="ql-code-block"> </dependency> </div> <div class="ql-code-block"> <dependency> </div> <div class="ql-code-block"> <groupId>com.alibaba</groupId> </div> <div class="ql-code-block"> <artifactId>fastjson</artifactId> </div> <div class="ql-code-block"> <version>1.2.83</version> </div> <div class="ql-code-block"> </dependency> </div> <div class="ql-code-block"> </dependencies> </div> </div> <p><br></p> <h4><strong style="color: rgb(52, 73, 94);">目录结构</strong></h4> <p><br></p> <p><span class="ql-font-monospace"><img src="https://pic.code-nav.cn/planet_post_image/1620773250336374786/2ma63xpd.jpeg"></span></p> <p><br></p> <h4><strong style="color: rgb(52, 73, 94);">具体实现</strong></h4> <p><br></p> <ol> <li data-list="ordered"><span class="ql-ui"></span>WebSocket需要一个配置文件,需要配置一个ServerEndpointExporter的Bean,ServerEndpointExporter的作用是自动扫描@ServerEndpoint所标记类,把该类注册成一个WebSocket连接类</li> </ol> <div class="ql-code-block-container"> <div class="ql-code-block"> /** </div> <div class="ql-code-block"> * WebSocket配置类 </div> <div class="ql-code-block"> */ </div> <div class="ql-code-block"> @Configuration </div> <div class="ql-code-block"> public class WebSocketConfig { </div> <div class="ql-code-block"> @Bean </div> <div class="ql-code-block"> public ServerEndpointExporter serverEndpointExporter(){ </div> <div class="ql-code-block"> return new ServerEndpointExporter(); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <ol> <li data-list="ordered"><span class="ql-ui"></span>具体的WebSocket服务端</li> </ol> <div class="ql-code-block-container"> <div class="ql-code-block"> /** </div> <div class="ql-code-block"> * WebSocket服务端,用于和前端连接 ServerEndpoint注解表示这个类是一个WebSocket连接类,括号里面的值表示客户端用于访问当前服务端的地址 </div> <div class="ql-code-block"> */ </div> <div class="ql-code-block"> @ServerEndpoint("/WebSocket/{userId}") </div> <div class="ql-code-block"> @Slf4j </div> <div class="ql-code-block"> @Component </div> <div class="ql-code-block"> public class WebSocket { </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> //记录连接的客户端 </div> <div class="ql-code-block"> private static final Map<String, Session> clients = new ConcurrentHashMap<String, Session>(); </div> <div class="ql-code-block"> //userId关联sid,解决一个userId连接多个服务端的问题 </div> <div class="ql-code-block"> private static final Map<String, Set<String>> connections = new ConcurrentHashMap<String, Set<String>>(); </div> <div class="ql-code-block"> //连接id </div> <div class="ql-code-block"> private String sid = null; </div> <div class="ql-code-block"> private String userId = null; </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> private static AtomicLong initialCount = new AtomicLong(0); </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> //判断是否连接 </div> <div class="ql-code-block"> public static boolean judgeConnect() { </div> <div class="ql-code-block"> if (CollectionUtils.isEmpty(clients.values())) { </div> <div class="ql-code-block"> log.info("未连接"); </div> <div class="ql-code-block"> return false; </div> <div class="ql-code-block"> } else { </div> <div class="ql-code-block"> log.info("已连接"); </div> <div class="ql-code-block"> return true; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> //给所有客户端发送消息 </div> <div class="ql-code-block"> public static void sendMessage(Message message) { </div> <div class="ql-code-block"> String messageStr = JSONObject.toJSONString(message); </div> <div class="ql-code-block"> for (Session value : clients.values()) { </div> <div class="ql-code-block"> try { </div> <div class="ql-code-block"> value.getBasicRemote().sendText(messageStr); </div> <div class="ql-code-block"> } catch (Exception e) { </div> <div class="ql-code-block"> log.error("发送消息错误", e); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> //给指定用户发送消息 </div> <div class="ql-code-block"> public static void sendMessageByUserId(Message message, String userId) { </div> <div class="ql-code-block"> if (!StringUtils.hasText(userId)) { </div> <div class="ql-code-block"> return; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> Set<String> clientSet = connections.get(userId); </div> <div class="ql-code-block"> if (CollectionUtils.isEmpty(clientSet)) { </div> <div class="ql-code-block"> return; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> String messageStr = JSONObject.toJSONString(message); </div> <div class="ql-code-block"> for (String sid : clientSet) { </div> <div class="ql-code-block"> Session session = clients.get(sid); </div> <div class="ql-code-block"> Optional.ofNullable(session).ifPresent(one -> { </div> <div class="ql-code-block"> try { </div> <div class="ql-code-block"> one.getBasicRemote().sendText(messageStr); </div> <div class="ql-code-block"> } catch (Exception e) { </div> <div class="ql-code-block"> log.error("发送消息错误", e); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> }); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> //连接成功 </div> <div class="ql-code-block"> @OnOpen </div> <div class="ql-code-block"> public void onOpen(Session session, @PathParam("userId") String userId) { </div> <div class="ql-code-block"> this.sid = UUID.randomUUID().toString(); </div> <div class="ql-code-block"> clients.put(this.sid, session); </div> <div class="ql-code-block"> this.userId = userId; </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> Set<String> clientSet = connections.get(userId); </div> <div class="ql-code-block"> if (CollectionUtils.isEmpty(clientSet)) { </div> <div class="ql-code-block"> clientSet = new HashSet<String>(); </div> <div class="ql-code-block"> connections.put(userId, clientSet); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> clientSet.add(this.sid); </div> <div class="ql-code-block"> log.info(this.sid + "已开启连接"); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> //连接关闭 </div> <div class="ql-code-block"> @OnClose </div> <div class="ql-code-block"> public void onClose() { </div> <div class="ql-code-block"> clients.remove(this.sid); </div> <div class="ql-code-block"> log.info(this.sid + "已断开连接"); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> //前端接受到消息的回调 </div> <div class="ql-code-block"> @OnMessage </div> <div class="ql-code-block"> public void onMessage(String message) { </div> <div class="ql-code-block"> log.info("前端已收到消息,返回消息为:" + message); </div> <div class="ql-code-block"> if ("消息已确认收到".equals(message)) { </div> <div class="ql-code-block"> initialCount.incrementAndGet(); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> //前端发生错误的回调 </div> <div class="ql-code-block"> @OnError </div> <div class="ql-code-block"> public void onError(Throwable error) { </div> <div class="ql-code-block"> log.error("发生了错误", error); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <ol> <li data-list="ordered"><span class="ql-ui"></span>Message类和他的子类OrderMessage,这里之所以抽象出一个Message父类,是为了保证扩展性,在WebSocket的SendMessage方法里可以接收任意的Message子类,达到发布任意的消息的功能</li> </ol> <div class="ql-code-block-container"> <div class="ql-code-block"> @Data </div> <div class="ql-code-block"> public abstract class Message { </div> <div class="ql-code-block"> private String title; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> /** </div> <div class="ql-code-block"> * 模拟支付信息 </div> <div class="ql-code-block"> */ </div> <div class="ql-code-block"> @EqualsAndHashCode(callSuper = true) </div> <div class="ql-code-block"> @Data </div> <div class="ql-code-block"> public class OrderMessage extends Message{ </div> <div class="ql-code-block"> private String status; </div> <div class="ql-code-block"> private BigDecimal cost; </div> <div class="ql-code-block"> } </div> </div> <ol> <li data-list="ordered"><span class="ql-ui"></span>OrderController用来模拟支付,在前后端连接建立之后,给订单设定“已支付”状态,然后通过WebSocket把消息发送给前端,前端接收后再具体执行剩下的操作</li> </ol> <div class="ql-code-block-container"> <div class="ql-code-block"> @RestController </div> <div class="ql-code-block"> @RequestMapping("/order") </div> <div class="ql-code-block"> public class OrderController { </div> <div class="ql-code-block"> @GetMapping("/payOrder/{orderId}") </div> <div class="ql-code-block"> public String payOrder(@PathVariable("orderId")String orderId){ </div> <div class="ql-code-block"> OrderMessage orderMessage = new OrderMessage(); </div> <div class="ql-code-block"> //模拟支付成功 </div> <div class="ql-code-block"> orderMessage.setTitle("订单id为"+orderId+"的订单已经被支付"); </div> <div class="ql-code-block"> orderMessage.setCost(BigDecimal.ONE); </div> <div class="ql-code-block"> orderMessage.setStatus("已支付"); </div> <div class="ql-code-block"> //发送消息 </div> <div class="ql-code-block"> WebSocket.sendMessage(orderMessage); </div> <div class="ql-code-block"> return "success"; </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> </div> <p><br></p> <h3><strong style="color: rgb(52, 73, 94);">前端的实现</strong></h3> <p><br></p> <p>前端的实现就通过一个简单的html代码模拟</p> <div class="ql-code-block-container"> <div class="ql-code-block"> <!DOCTYPE html> </div> <div class="ql-code-block"> <html lang="en"> </div> <div class="ql-code-block"> <head> </div> <div class="ql-code-block"> <meta charset="UTF-8"> </div> <div class="ql-code-block"> <title>SseEmitter</title> </div> <div class="ql-code-block"> </head> </div> <div class="ql-code-block"> <body> </div> <div class="ql-code-block"> <div id="message"></div> </div> <div class="ql-code-block"> </body> </div> <div class="ql-code-block"> <script> </div> <div class="ql-code-block"> var limitConnect = 0; </div> <div class="ql-code-block"> // 如果订单状态为未支付,就建立连接,如果订单超时或者已经支付,就不建立连接,现在默认订单未支付 </div> <div class="ql-code-block"> // 初始化,建立连接 </div> <div class="ql-code-block"> init(); </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> function init() { </div> <div class="ql-code-block"> // 8080未默认端口,可自行替换,这里的路径和后端的@ServerEndpoint的路径需要对应上 </div> <div class="ql-code-block"> var ws = new WebSocket('ws://localhost:8080/WebSocket/1'); </div> <div class="ql-code-block"> // 获取连接状态 </div> <div class="ql-code-block"> console.log('WebSocket连接状态:' + ws.readyState); </div> <div class="ql-code-block"> //监听是否连接成功 </div> <div class="ql-code-block"> ws.onopen = function () { </div> <div class="ql-code-block"> console.log('WebSocket连接状态:' + ws.readyState); </div> <div class="ql-code-block"> limitConnect = 0; </div> <div class="ql-code-block"> //连接成功则发送一个数据 </div> <div class="ql-code-block"> ws.send('我们建立连接啦'); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> // 接听服务器发回的信息并处理展示 </div> <div class="ql-code-block"> ws.onmessage = function (data) { </div> <div class="ql-code-block"> console.log('接收到来自服务器的消息:'); </div> <div class="ql-code-block"> console.log(data); </div> <div class="ql-code-block"> //发起消息回调,告诉后端,前端已经收到消息 </div> <div class="ql-code-block"> ws.send("前端已接收到消息") </div> <div class="ql-code-block"> //收到消息之后,代表订单已被支付,可以刷新页面,或者跳转到支付成功的页面等其他操作 </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> // 监听连接关闭事件 </div> <div class="ql-code-block"> ws.onclose = function () { </div> <div class="ql-code-block"> // 监听整个过程中websocket的状态 </div> <div class="ql-code-block"> console.log('WebSocket连接状态:' + ws.readyState); </div> <div class="ql-code-block"> reconnect(); </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> // 监听并处理error事件 </div> <div class="ql-code-block"> ws.onerror = function (error) { </div> <div class="ql-code-block"> console.log(error); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </div> <div class="ql-code-block"> function reconnect() { </div> <div class="ql-code-block"> limitConnect++; </div> <div class="ql-code-block"> console.log("重连第" + limitConnect + "次"); </div> <div class="ql-code-block"> setTimeout(function () { </div> <div class="ql-code-block"> init(); </div> <div class="ql-code-block"> }, 2000); </div> <div class="ql-code-block"> } </div> <div class="ql-code-block"> </script> </div> <div class="ql-code-block"> </html> </div> </div> <p><br></p> <h2><strong style="color: rgb(52, 73, 94);">实测</strong></h2> <p><br></p> <ol> <li data-list="ordered"><span class="ql-ui"></span>启动后端项目,然后打开上述的html文件,我们可以在控制台看到</li> <li data-list="ordered"><span class="ql-ui"></span><span class="ql-font-monospace"><img src="https://pic.code-nav.cn/planet_post_image/1620773250336374786/nzcai61k.jpeg"></span></li> </ol> <p>前端的控制台也可以看到</p> <p><span class="ql-font-monospace"><img src="https://pic.code-nav.cn/planet_post_image/1620773250336374786/s58mkys3.jpeg"></span></p> <p>有两个连接状态打印对应着前端的代码,未连接前会打印一次,连接之后又会打印一次,并且向后端发送消 息,说明前后端已经连接</p> <p><span class="ql-font-monospace"><img src="https://pic.code-nav.cn/planet_post_image/1620773250336374786/1384wtrf.jpeg"></span></p> <ol> <li data-list="ordered"><span class="ql-ui"></span>去调用OrderController里面的payOrder方法模拟一次支付,可见请求成功</li> <li data-list="ordered"><span class="ql-ui"></span><span class="ql-font-monospace"><img src="https://pic.code-nav.cn/planet_post_image/1620773250336374786/nhk9w1qz.jpeg"></span></li> <li data-list="ordered"><span class="ql-ui"></span>查看前端控制台,发现后端给前端发送的消息已经打印在控制台了</li> <li data-list="ordered"><span class="ql-ui"></span><span class="ql-font-monospace"><img src="https://pic.code-nav.cn/planet_post_image/1620773250336374786/sth4xvp4.jpeg"></span></li> </ol> <p>然后看后端控制台,发现后端也接收到前端确认消息的信息了</p> <p><span class="ql-font-monospace"><img src="https://pic.code-nav.cn/planet_post_image/1620773250336374786/x7rkyc48.jpeg"></span></p> <ol> <li data-list="ordered"><span class="ql-ui"></span>这样就说明前后端的连接已经建立成功,并且能实时进行通讯了</li> </ol> <p><br></p> <h2><strong style="color: rgb(52, 73, 94);">总结</strong></h2> <p><br></p> <p>WebSocket的使用总结起来就一个词,方便,我可以在上述demo的基础上实现多种场景的开发,比如一开始提到的核验之后进行页面刷新,订单超时之后跳转页面等,然后后端其实还有更加强大的网络编程框架Netty,虽然我没用过,但是应该也有类似的功能来实现这种前后端长连接的建立</p> <p><br></p> </div> </body> </html>
