侧边栏壁纸
博主头像
帆云素材 博主等级

苟利国家生死以,岂因祸福避趋之

  • 累计撰写 26 篇文章
  • 累计创建 26 个标签
  • 累计收到 2 条评论

目 录CONTENT

文章目录

流式传输-SSE

智慧的格子衫
2025-01-07 / 0 评论 / 0 点赞 / 15 阅读 / 0 字

随着Web应用的不断发展,实时数据传输的需求变得越来越普遍。传统的轮询方法不仅效率低下,而且在高并发情况下会对服务器造成不必要的压力。为了解决这个问题,Server-Sent Events (SSE) 应运而生,它允许服务器端主动向客户端推送更新。

Server-Sent Events

Server-Sent Events 是一种允许服务器向浏览器发送实时更新的技术。不同于WebSocket的全双工通信方式,SSE更专注于单向的数据流,即从服务器到客户端的数据推送。

这种方式对于需要实时更新的场景非常有用,当前主流的大模型平台,比如ChatGPT、通义千问、文心一言,对话时采用的就是SSE。

SSE 本质是一个基于 http 协议的通信技术。

SSE的应用

引入依赖

spring-boot-starter-web 中默认已经引用了 sse,所以我们不需要额外引入其他依赖

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
</dependency>

基本用法

@RestController
public class SseController {

    @GetMapping("/sse")
    public SseEmitter handleSse() {
        SseEmitter emitter = new SseEmitter();

        // 异步处理发送事件
        new Thread(() -> {
            try {
                // 推送事件
                emitter.send("实时消息:你好!");
                Thread.sleep(1000);  // 模拟延时
                emitter.send("实时消息:更新来了!");
                emitter.complete(); // 结束推送
            } catch (Exception e) {
                emitter.completeWithError(e); // 异常处理
            }
        }).start();

        return emitter;
    }
}

客户端代码

const eventSource = new EventSource("/sse");

// 处理服务器推送的消息
eventSource.onmessage = function(event) {
    console.log(event.data);
};

// 处理连接关闭或错误
eventSource.onerror = function() {
    console.log("连接出现问题,自动重连...");
};

效果预览

msg

为什么会这么多,那是因为客户端的自动重连机制,无需我们手动维护,客户端会自动发起重连。

SseEmitter

其实sse的核心,就是SseEmitter这个类,是 Spring 提供的一个类,用于处理 Server-Sent Events (SSE)。它允许服务器端以流的形式推送事件给客户端,而不需要客户端不断轮询服务器。

  1. 构造函数
  • SseEmitter():创建一个默认超时时间的 SseEmitter 实例。默认超时为 30 秒。
  • SseEmitter(Long timeout):创建一个带有自定义超时时间的 SseEmitter 实例。
    • timeout:指定以毫秒为单位的超时时间。如果设置为 0L,则连接永远不会超时。
  1. 核心方法
  • send(Object object):向客户端发送一条消息。
    • object:要发送的数据,可以是任何类型的对象。
    • 此方法会将数据直接发送到客户端,并在响应体中流式返回。
  • send(SseEmitter.SseEventBuilder event):以事件构建器的形式发送一条消息。
    • SseEventBuilder 是用来构建发送事件的一个内部类,允许你自定义事件的各个属性,如 iddataname 等。
  • complete():表示 SSE 流完成并关闭连接。服务器告诉客户端,事件流已经结束。
  • completeWithError(Throwable ex):在发生错误时关闭连接,并以错误的形式告知客户端。
  1. 回调函数
  • onCompletion(Runnable callback):指定当 SSE 连接完成(正常关闭)时执行的回调函数。
  • onTimeout(Runnable callback):指定当连接超时时执行的回调函数。
  • onError(Consumer<Throwable> callback):指定当发生错误时执行的回调函数。这个错误可能是由于网络连接问题、客户端断开等原因。
  1. SseEventBuilder 内部类

SseEmitter.SseEventBuilder 是用于构建 SSE 事件的一个类,允许自定义事件的各个部分:

  • SseEmitter.event():返回一个新的 SseEventBuilder 实例,用于构建事件。
  • id(String id):设置事件的唯一标识符。客户端可以通过这个 ID 识别和处理事件。
  • name(String name):设置事件的名称。客户端可以通过这个名称识别不同类型的事件,SSE 响应中会显示为 event: name
  • data(Object data):设置要发送的数据。可以是文本、JSON 等类型,最终会在客户端的 SSE 流中显示为 data: xxx
  • reconnectTime(long milliseconds):告诉客户端在多少毫秒后尝试重新连接。如果连接中断,客户端会在指定时间后自动重连。
  • comment(String comment):向客户端发送一条注释(不会触发事件)。

目前 SseEmitter 是基于 每个客户端请求独立管理 的对象,因此不适合将其直接交由 Spring 管理为单例或共享对象。

每次请求应手动创建新的 SseEmitter 实例,并配置合适的超时时间。对于每个连接,SseEmitter 都是短暂的,使用完毕后应该调用 complete()completeWithError() 方法来释放资源。

SSE 与 WebSocket 对比

lp6iz1

提供一个sse的工具类

import lombok.extern.slf4j.Slf4j;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;

import java.io.IOException;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

/**
 * Sse 客户端
 *
 * @author fyf
 * @since 2024/9/23 10:20
 */
@Slf4j
public class SseClient {

    private static final Map<String, SseEmitter> sseEmitterMap = new ConcurrentHashMap<>();

    /**
     * 10分钟超时
     */
    private static final long TIMEOUT = 600000L;

    /**
     * 客户端重连时间
     */
    private static final long RECONNECT_TIME = 60000L;
    
    /**
     * 创建连接
     *
     * @param uid sse连接id(用户id)
     * @return SseEmitter 对象
     */
    public static SseEmitter createSse(String uid) {
        if (sseEmitterMap.containsKey(uid)) {
            log.warn("[{}]连接已存在", uid);
            return sseEmitterMap.get(uid);
        }
        SseEmitter sseEmitter = getSseEmitter(uid);
        try {
            // 发送初始化消息,设置重连时间
            sseEmitter.send(SseEmitter.event().reconnectTime(RECONNECT_TIME));
        } catch (IOException e) {
            log.error("[{}]创建连接失败: {}", uid, e.getMessage());
        }
        sseEmitterMap.put(uid, sseEmitter);
        log.info("[{}]创建SSE连接成功", uid);
        return sseEmitter;
    }


    /**
     * 给指定用户发送消息-发送失败会创建新的连接并重试发送
     *
     * @param uid 用户id
     * @param message 消息内容
     * @return 发送成功返回true,失败返回false
     */
    public static boolean sendMessage(String uid, String message) {
        return sendMessage(uid, message, true);
    }

    /**
     * 给指定用户发送消息
     *
     * @param uid 用户id
     * @param message 消息内容
     * @param retry 是否创建连接重试
     * @return 发送成功返回true,失败返回false
     */
    public static boolean sendMessage(String uid, String message, boolean retry) {
        if (message == null || message.isEmpty()) {
            log.warn("发送失败,消息为空");
            return false;
        }

        SseEmitter sseEmitter = sseEmitterMap.get(uid);
        if (sseEmitter == null) {
            log.warn("消息推送失败,用户[{}]没有连接", uid);
            if (!retry) {
                return false;
            }
            log.info("[{}]尝试重新创建连接...", uid);
            sseEmitter = createSse(uid);
        }

        try {
            sseEmitter.send(SseEmitter.event().id(uid).data(message));
            log.debug("[{}]消息发送成功,消息内容: {}", uid, message);
            sseEmitter.complete();
            closeSse(uid);
            return true;
        } catch (IOException e) {
            log.error("[{}]消息发送失败,连接中断: {}", uid, e.getMessage());
            // 关闭失效的连接
            closeSse(uid);
            if (!retry) {
                return false;
            }
            // 尝试重新创建连接并再次发送消息
            log.info("[{}]尝试重新创建连接并重新发送消息...", uid);
            sseEmitter = createSse(uid);
            try {
                sseEmitter.send(SseEmitter.event().id(uid).data(message));
                log.info("[{}]重新发送消息成功,消息内容: {}", uid, message);
                return true;
            } catch (IOException ex) {
                log.error("[{}]重新发送消息失败: {}", uid, ex.getMessage());
                return false;
            }
        }
    }

    /**
     * 手动关闭连接
     *
     * @param uid 用户id
     */
    public static void closeSse(String uid) {
        SseEmitter sseEmitter = sseEmitterMap.get(uid);
        if (sseEmitter != null) {
            sseEmitter.complete(); // 结束连接
            sseEmitterMap.remove(uid);
            log.info("[{}]连接已手动关闭", uid);
        } else {
            log.warn("[{}]连接已不存在", uid);
        }
    }

    private static SseEmitter getSseEmitter(String uid) {
        SseEmitter sseEmitter = new SseEmitter(TIMEOUT);
        // 完成后回调,清理资源
        sseEmitter.onCompletion(() -> {
            log.info("[{}]结束连接", uid);
            sseEmitterMap.remove(uid);
        });
        // 超时回调
        sseEmitter.onTimeout(() -> {
            log.info("[{}]连接超时", uid);
            sseEmitterMap.remove(uid);
        });
        // 异常回调,发生异常时关闭连接并清理
        sseEmitter.onError(throwable -> {
            log.error("[{}]连接异常: {}", uid, throwable.getMessage());
            closeSse(uid);
        });
        return sseEmitter;
    }
}

修改nginx

如果使用nginx代理,还需要禁用nginx的缓存

 location / {
    proxy_pass http://127.0.0.1:8100;
    proxy_set_header Host $http_host;
    proxy_set_header X-Real-IP $remote_addr;
    proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
    proxy_set_header Connection '';
    proxy_http_version 1.1;  # 重要:确保使用HTTP/1.1协议
    proxy_set_header Upgrade $http_upgrade;
    proxy_set_header Connection 'upgrade';

    # 添加以下配置以处理SSE
    proxy_buffering off;
    proxy_cache off;
}

小结

  • 相比轮询,SSE 通过长连接减少了网络开销和服务器压力。
  • SseEmitter 适用于需要服务器实时推送数据的场景,特别是实时通知、动态更新等需求
  • 相比websocket,在某些场景下,SSE 更加易用,且占用的资源较少
0

评论区