Web服务端通过SSE推送消息给浏览器客户端的实现方案(附详细代码和仓库地址)

本文主要是介绍Web服务端通过SSE推送消息给浏览器客户端的实现方案(附详细代码和仓库地址),希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

目录

    • 1、SSE(Server-Sent Events)简介
    • 2、SSE 的工作原理
    • 3、SSE 与客户端轮询的区别和优势比较
      • 区别
      • 优势
    • 4、SSE简单实现(单机应用Demo)
      • 演示效果
      • SSE-Demo仓库地址
      • 下面直接贴代码:
      • 前端实现:
      • 后端实现:
    • 5、SSE简单实现(分布式应用Demo)
      • SSE-Demo仓库地址
      • 关键代码
      • 方案说明

1、SSE(Server-Sent Events)简介

Server-Sent Events (SSE) 是一种基于 HTTP 协议的服务器推送技术,允许服务器通过单个持久连接向客户端发送实时更新。客户端使用标准的 EventSource API 来接收服务器推送的事件,这种通信方式非常适合实时应用,如消息通知、股票行情更新、社交媒体更新等。

2、SSE 的工作原理

  • 单向连接:SSE 建立的是单向通道,即服务器向客户端推送数据,客户端只能接收,不能发送。
  • 持久连接:SSE 使用的是长连接(Long Polling),即连接一旦建立,将会持续存在,直到客户端或服务器关闭连接。
  • 文本数据:SSE 通过 text/event-stream MIME 类型传输数据,数据是纯文本格式。
  • 自动重连:如果连接中断,EventSource 会自动尝试重新连接,确保客户端能够接收后续的推送。

3、SSE 与客户端轮询的区别和优势比较

客户端轮询(Client Polling) 是一种传统的客户端从服务器请求数据的方式。客户端会定期向服务器发送请求,检查是否有新数据可用。

区别

连接方式:
SSE:建立后服务器主动推送数据,连接是持久的,数据在有更新时实时传递。
轮询:客户端定期发送请求获取数据,连接是间歇性的。

实时性:
SSE:数据几乎是实时推送的,延迟极低。
轮询:数据获取延迟取决于轮询的频率,频率高则延迟低,但频率低可能导致数据延迟。

网络和服务器负载:
SSE:由于是单个持久连接,减少了频繁的请求与响应开销,降低了服务器负载。
轮询:频繁的请求会增加服务器和网络的负担,尤其是在轮询频率较高时。

连接控制:
SSE:自动处理连接中断和重连,客户端实现简单。
轮询:需要客户端定期发起请求,且如果请求频率不当,可能导致资源浪费。

数据传输效率:
SSE:只在有数据更新时推送,传输效率高。
轮询:即使没有数据更新,客户端也会定期请求,效率低下。

优势

SSE 的优势:
更高效的网络和服务器资源利用率。
实时性更高,延迟更低。
实现简单,特别是在浏览器环境中,支持自动重连和事件处理。
适合需要频繁更新但客户端无需响应的场景。

客户端轮询的优势:
在不支持 SSE 的环境下仍然可以使用。
实现和理解相对简单,兼容性更好。

4、SSE简单实现(单机应用Demo)

演示效果

在这里插入图片描述

SSE-Demo仓库地址

https://github.com/deepjava-gm/SSE-Demo.git

下面直接贴代码:

前端实现:

<!DOCTYPE html>
<html lang="en">
<head><meta charset="UTF-8"><meta name="viewport" content="width=device-width, initial-scale=1.0"><title>SSE 用户消息推送 Demo</title><style>body {font-family: Arial, sans-serif;margin: 20px;}h1 {text-align: center;}.container {max-width: 600px;margin: 0 auto;padding: 20px;border: 1px solid #ccc;border-radius: 8px;}.form-group {margin-bottom: 15px;}label {display: block;margin-bottom: 5px;font-weight: bold;}input[type="text"] {width: 100%;padding: 8px;box-sizing: border-box;border: 1px solid #ccc;border-radius: 4px;}button {padding: 10px 20px;color: white;background-color: #007bff;border: none;border-radius: 4px;cursor: pointer;}button:hover {background-color: #0056b3;}#messages {margin-top: 20px;}.message {background-color: #f1f1f1;padding: 10px;margin-bottom: 10px;border-radius: 4px;}.status {margin-top: 10px;font-weight: bold;}.success {color: green;}.error {color: red;}</style>
</head>
<body><h1>SSE 用户消息推送 Demo</h1><div class="container"><!-- 连接部分 --><div class="form-group"><label for="userId">用户 ID:</label><input type="text" id="userId" placeholder="请输入您的用户 ID"><br><br><button onclick="connect()">连接</button></div><div id="connectionStatus" class="status"></div><!-- 消息发送部分 --><div class="form-group"><label for="targetUserId">目标用户 ID:</label><input type="text" id="targetUserId" placeholder="请输入目标用户 ID"></div><div class="form-group"><label for="message">消息内容:</label><input type="text" id="message" placeholder="请输入要发送的消息"></div><button onclick="sendMessage()">推送消息</button><!-- 消息显示部分 --><div id="messages"></div></div><script>let eventSource;let currentUserId = '';function connect() {const userId = document.getElementById('userId').value;const connectionStatus = document.getElementById('connectionStatus');if (!userId) {alert('请输入用户 ID');return;}// 显示连接状态为“连接中”connectionStatus.textContent = '已连接...';connectionStatus.className = 'status';eventSource = new EventSource(`http://localhost:9999/sse/connect/${userId}`);currentUserId = userId; // 保存当前用户 IDeventSource.onopen = function() {connectionStatus.textContent = '接收成功';connectionStatus.className = 'status success';};eventSource.onmessage = function(event) {try {// 解析 JSON 消息const data = JSON.parse(event.data);const newElement = document.createElement('div');newElement.className = 'message';newElement.innerText = `用户 ${data.senderId} 接收的消息: ${data.message}`;document.getElementById('messages').appendChild(newElement);} catch (e) {console.error('消息解析错误:', e);}};eventSource.onerror = function(event) {connectionStatus.textContent = '连接失败,请检查网络或服务器';connectionStatus.className = 'status error';console.error("连接错误: ", event);eventSource.close();};}function sendMessage() {const targetUserId = document.getElementById('targetUserId').value;const message = document.getElementById('message').value;if (!targetUserId || !message) {alert('请填写目标用户 ID 和消息内容');return;}// 发送 GET 请求推送消息fetch(`http://localhost:9999/sse/push/${targetUserId}?message=${encodeURIComponent(message)}`, {method: 'GET'}).then(response => {if (response.ok) {console.log('消息发送成功');} else {console.log('消息发送失败');}}).catch(error => {console.error('发送错误:', error);alert('消息发送失败');});}</script>
</body>
</html>

后端实现:

启动类:

package io.github.deepjava;import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;import java.util.concurrent.ConcurrentHashMap;@SpringBootApplication
public class SseDemoApplication {public static void main(String[] args) {SpringApplication.run(SseDemoApplication.class);}// 注入一个全局缓存 用来保存不同用户的SSE连接信息@Bean("userSSEMap")public ConcurrentHashMap<String, SseEmitter> getUserSSEMap(){return new ConcurrentHashMap<>();}}

Controller:

package io.github.deepjava.controller;import io.github.deepjava.service.SseService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;import javax.annotation.Resource;
import java.util.concurrent.ConcurrentHashMap;@RestController
@RequestMapping("/sse")
@CrossOrigin(origins = "*")
@Slf4j
public class SseController {@Resource(name = "userSSEMap")private ConcurrentHashMap<String, SseEmitter> userSSEMap;@Resourceprivate SseService sseService;// 连接方法:为用户 ID 注册 SSE 链接@GetMapping("/connect/{userId}")public SseEmitter connect(@PathVariable String userId) {SseEmitter emitter = new SseEmitter(0L); // 设置超时时间为无限大userSSEMap.put(userId, emitter);// 连接正常关闭回调 移除连接emitter.onCompletion(() -> {userSSEMap.remove(userId);log.info("连接正常关闭回调 移除连接");});// 连接超时回调 移除连接emitter.onTimeout(() -> {userSSEMap.remove(userId);log.info("连接超时回调 移除连接");});// 连接出错回调 移除连接emitter.onError((e) -> {userSSEMap.remove(userId);log.info("连接出错回调 移除连接");});log.info("连接成功!");return emitter;}// 推送方法:根据用户 ID 发送消息@GetMapping("/push/{userId}")public void push(@PathVariable String userId, @RequestParam String message) {sseService.extracted(userId, message);}}

Service:

package io.github.deepjava.service;import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;import javax.annotation.Resource;
import java.io.IOException;
import java.util.concurrent.ConcurrentHashMap;@Service
@Slf4j
public class SseService {@Resource(name = "userSSEMap")private ConcurrentHashMap<String, SseEmitter> clients;public void extracted(String userId, String message) {SseEmitter emitter = clients.get(userId);if (emitter != null) {try {// 创建包含用户 ID 和消息内容的 JSON 对象String jsonMessage = String.format("{\"senderId\":\"%s\", \"message\":\"%s\"}", userId, message);emitter.send(jsonMessage);log.info("消息推送成功!");} catch (IOException e) {clients.remove(userId);log.info("消息推送失败!");}}}}

配置文件:application.properties

spring.application.name=sse-demo
server.port=9999

Maven的pom文件:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"><modelVersion>4.0.0</modelVersion><parent><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-parent</artifactId><version>2.7.1</version></parent><groupId>org.example</groupId><artifactId>SSE-Demo</artifactId><version>1.0-SNAPSHOT</version><dependencies><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-web</artifactId></dependency><!--lombok  --><dependency><groupId>org.projectlombok</groupId><artifactId>lombok</artifactId></dependency></dependencies><properties><maven.compiler.source>8</maven.compiler.source><maven.compiler.target>8</maven.compiler.target></properties></project>

5、SSE简单实现(分布式应用Demo)

注意:
SSE 连接(如 SseEmitter)是持久化的、与具体服务器实例相关联的动态对象,无法直接存储在 Redis 等外部存储中。Redis 主要用于消息传递和共享数据,但无法直接管理活跃的连接。所以下面方案仅使用 Redis 进行消息广播。

解决方案概述
为了在分布式环境中实现 SSE,通常采用以下架构:

  • 每个服务器实例维护本地的 SSE 连接:每个实例只管理与自身连接的客户端。
  • 使用 Redis 进行消息广播:当需要向特定用户推送消息时,将消息发布到 Redis 频道。所有实例订阅该频道,并检查自己是否有需要向某个用户推送的连接。
  • 用户与实例的映射:使用 Redis 存储用户与服务器实例的映射信息,确保消息能够被正确路由到处理该用户连接的实例。
    虽然无法完全将连接信息存储在 Redis 中,但通过这种方式,可以有效地在分布式环境中管理 SSE 连接和消息推送。

这里只贴主要的后端代码:完整代码去下载仓库代码看。

SSE-Demo仓库地址

https://github.com/deepjava-gm/SSE-Demo.git

关键代码

redis配置:

 // 配置redis的序列化@Beanpublic RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {RedisTemplate<String, Object> template = new RedisTemplate<>();template.setConnectionFactory(factory);template.setKeySerializer(new StringRedisSerializer());template.setValueSerializer(new StringRedisSerializer());return template;}@Beanpublic RedisMessageListenerContainer redisMessageListenerContainer(RedisConnectionFactory connectionFactory) {RedisMessageListenerContainer container = new RedisMessageListenerContainer();container.setConnectionFactory(connectionFactory);return container;}

redis监听主题:

package io.github.deepjava.controller;import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.connection.MessageListener;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.data.redis.listener.ChannelTopic;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
import org.springframework.data.redis.listener.adapter.MessageListenerAdapter;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;import javax.annotation.PostConstruct;
import javax.annotation.Resource;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.ConcurrentHashMap;@RestController
@RequestMapping("/dis/sse")
@CrossOrigin(origins = "*")
@Slf4j
public class DistributedSseController {@Resource(name = "userSSEMap")private ConcurrentHashMap<String, SseEmitter> userSSEMap;@Resourceprivate RedisTemplate<String, String> redisTemplate;@Resourceprivate RedisMessageListenerContainer redisMessageListenerContainer;private final ChannelTopic topic = new ChannelTopic("sse-messages");@PostConstructpublic void init() {// 订阅 Redis 频道redisMessageListenerContainer.addMessageListener(new MessageListenerAdapter((MessageListener) (message, pattern) -> {String payload = new String(message.getBody(), StandardCharsets.UTF_8);// 假设消息的格式为 "userId:message"String[] parts = payload.split(":", 2);if (parts.length == 2) {String userId = parts[0];String userMessage = parts[1];// 发送消息给本地的 SSE 连接SseEmitter emitter = userSSEMap.get(userId);if (emitter != null) {try {String jsonMessage = String.format("{\"senderId\":\"%s\", \"message\":\"%s\"}", userId, userMessage);emitter.send(jsonMessage);} catch (IOException e) {emitter.completeWithError(e);userSSEMap.remove(userId);}}}}), topic);}// 连接方法:为用户 ID 注册 SSE 链接@GetMapping("/connect/{userId}")public SseEmitter connect(@PathVariable String userId) {SseEmitter emitter = new SseEmitter(0L); // 设置超时时间为无限大userSSEMap.put(userId, emitter);// 连接正常关闭回调 移除连接emitter.onCompletion(() -> {userSSEMap.remove(userId);log.info("连接正常关闭回调 移除连接");});// 连接超时回调 移除连接emitter.onTimeout(() -> {userSSEMap.remove(userId);log.info("连接超时回调 移除连接");});// 连接出错回调 移除连接emitter.onError((e) -> {userSSEMap.remove(userId);log.info("连接出错回调 移除连接");});log.info("连接成功!");return emitter;}@GetMapping("/push/{userId}")public void push(@PathVariable String userId, String message) {// 将消息发布到 Redis 频道redisTemplate.convertAndSend(topic.getTopic(), userId + ":" + message);}
}

方案说明

SSE 连接管理:
使用 ConcurrentHashMap<String, SseEmitter> 存储用户的连接信息,每个服务器实例只维护与自身连接的客户端。
connect 方法用于创建 SSE 连接并保存到本地缓存。

Redis 消息广播:
通过 Redis 的 发布订阅(Pub/Sub) 机制,所有实例订阅同一个频道(sse-messages)。
push 方法将消息发布到 Redis 频道,所有订阅了该频道的实例都会收到消息,并检查是否有对应的连接需要推送。

这篇关于Web服务端通过SSE推送消息给浏览器客户端的实现方案(附详细代码和仓库地址)的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1128040

相关文章

python生成随机唯一id的几种实现方法

《python生成随机唯一id的几种实现方法》在Python中生成随机唯一ID有多种方法,根据不同的需求场景可以选择最适合的方案,文中通过示例代码介绍的非常详细,需要的朋友们下面随着小编来一起学习学习... 目录方法 1:使用 UUID 模块(推荐)方法 2:使用 Secrets 模块(安全敏感场景)方法

Windows环境下解决Matplotlib中文字体显示问题的详细教程

《Windows环境下解决Matplotlib中文字体显示问题的详细教程》本文详细介绍了在Windows下解决Matplotlib中文显示问题的方法,包括安装字体、更新缓存、配置文件设置及编码調整,并... 目录引言问题分析解决方案详解1. 检查系统已安装字体2. 手动添加中文字体(以SimHei为例)步骤

MySQL 迁移至 Doris 最佳实践方案(最新整理)

《MySQL迁移至Doris最佳实践方案(最新整理)》本文将深入剖析三种经过实践验证的MySQL迁移至Doris的最佳方案,涵盖全量迁移、增量同步、混合迁移以及基于CDC(ChangeData... 目录一、China编程JDBC Catalog 联邦查询方案(适合跨库实时查询)1. 方案概述2. 环境要求3.

Spring StateMachine实现状态机使用示例详解

《SpringStateMachine实现状态机使用示例详解》本文介绍SpringStateMachine实现状态机的步骤,包括依赖导入、枚举定义、状态转移规则配置、上下文管理及服务调用示例,重点解... 目录什么是状态机使用示例什么是状态机状态机是计算机科学中的​​核心建模工具​​,用于描述对象在其生命

Spring Boot 结合 WxJava 实现文章上传微信公众号草稿箱与群发

《SpringBoot结合WxJava实现文章上传微信公众号草稿箱与群发》本文将详细介绍如何使用SpringBoot框架结合WxJava开发工具包,实现文章上传到微信公众号草稿箱以及群发功能,... 目录一、项目环境准备1.1 开发环境1.2 微信公众号准备二、Spring Boot 项目搭建2.1 创建

IntelliJ IDEA2025创建SpringBoot项目的实现步骤

《IntelliJIDEA2025创建SpringBoot项目的实现步骤》本文主要介绍了IntelliJIDEA2025创建SpringBoot项目的实现步骤,文中通过示例代码介绍的非常详细,对大家... 目录一、创建 Spring Boot 项目1. 新建项目2. 基础配置3. 选择依赖4. 生成项目5.

nginx -t、nginx -s stop 和 nginx -s reload 命令的详细解析(结合应用场景)

《nginx-t、nginx-sstop和nginx-sreload命令的详细解析(结合应用场景)》本文解析Nginx的-t、-sstop、-sreload命令,分别用于配置语法检... 以下是关于 nginx -t、nginx -s stop 和 nginx -s reload 命令的详细解析,结合实际应

Spring boot整合dubbo+zookeeper的详细过程

《Springboot整合dubbo+zookeeper的详细过程》本文讲解SpringBoot整合Dubbo与Zookeeper实现API、Provider、Consumer模式,包含依赖配置、... 目录Spring boot整合dubbo+zookeeper1.创建父工程2.父工程引入依赖3.创建ap

SpringBoot3.X 整合 MinIO 存储原生方案

《SpringBoot3.X整合MinIO存储原生方案》本文详细介绍了SpringBoot3.X整合MinIO的原生方案,从环境搭建到核心功能实现,涵盖了文件上传、下载、删除等常用操作,并补充了... 目录SpringBoot3.X整合MinIO存储原生方案:从环境搭建到实战开发一、前言:为什么选择MinI

Linux下删除乱码文件和目录的实现方式

《Linux下删除乱码文件和目录的实现方式》:本文主要介绍Linux下删除乱码文件和目录的实现方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录linux下删除乱码文件和目录方法1方法2总结Linux下删除乱码文件和目录方法1使用ls -i命令找到文件或目录