You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

能否在响应过滤器中使用ServiceStack SSE?如何初始化连接及管理实例?

SSE相关问题解答

能否在响应过滤器中发送SSE消息?

可以,但需要注意几个核心要点:

  • 必须在响应未提交时操作,确保能修改响应头并获取响应流。
  • 要正确设置SSE必备响应头:Content-Type: text/event-stream、Cache-Control: no-cache、Connection: keep-alive,避免浏览器或代理缓存响应或提前断开连接。
  • 过滤器执行完成后不要立即关闭响应流,需保持连接打开以持续发送事件。

示例代码(以ServiceStack响应过滤器为例):

public class SseResponseFilter : IResponseFilter
{
    public void Execute(IHttpRequest req, IHttpResponse res, object dto)
    {
        // 仅处理指定的SSE请求路径
        if (!req.PathInfo.Equals("/sse", StringComparison.OrdinalIgnoreCase))
            return;

        // 设置SSE响应头
        res.ContentType = "text/event-stream";
        res.Headers.Add("Cache-Control", "no-cache");
        res.Headers.Add("Connection", "keep-alive");
        res.Headers.Add("Access-Control-Allow-Origin", "*"); // 跨域场景需添加

        // 获取响应流并发送初始连接确认
        var writer = new StreamWriter(res.OutputStream) { AutoFlush = true };
        writer.WriteLine("data: 连接已建立\n");

        // 将连接存入线程安全集合,用于后续推送
        var clientId = req.GetSessionId() ?? Guid.NewGuid().ToString();
        SseConnectionManager.AddConnection(clientId, writer);

        // 监听请求取消事件,及时清理连接
        req.RequestAborted.Register(() => 
            SseConnectionManager.RemoveConnection(clientId));
    }
}

初始化、连接和发布SSE事件的步骤

初始化与连接

  1. 服务器端初始化

    • 注册响应过滤器(或中间件),识别SSE请求并设置响应头。
    • 维护一个线程安全集合(如ConcurrentDictionary<string, StreamWriter>)存储活跃连接,关联客户端唯一标识(如SessionId、自定义ID)。
    • 监听请求取消/断开事件,及时清理无效连接。
  2. 客户端连接
    使用浏览器原生EventSource API建立连接:

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

eventSource.onopen = () => console.log('SSE连接已打开');
eventSource.onmessage = (e) => console.log('收到消息:', e.data);
eventSource.onerror = (err) => console.error('SSE连接出错:', err);

发布SSE事件

按照SSE规范格式写入响应流,支持带事件类型、ID等字段:

public static class SseConnectionManager
{
    private static readonly ConcurrentDictionary<string, StreamWriter> _activeConnections = new();

    public static void AddConnection(string clientId, StreamWriter writer) =>
        _activeConnections.TryAdd(clientId, writer);

    public static void RemoveConnection(string clientId)
    {
        if (_activeConnections.TryRemove(clientId, out var writer))
        {
            writer.Dispose();
        }
    }

    // 发送普通消息
    public static void SendMessage(string clientId, string message)
    {
        if (_activeConnections.TryGetValue(clientId, out var writer))
        {
            try
            {
                writer.WriteLine($"data: {message}\n");
            }
            catch (IOException)
            {
                // 连接已断开,自动清理
                RemoveConnection(clientId);
            }
        }
    }

    // 发送带事件类型的消息
    public static void SendEvent(string clientId, string eventType, string message)
    {
        if (_activeConnections.TryGetValue(clientId, out var writer))
        {
            try
            {
                writer.WriteLine($"event: {eventType}");
                writer.WriteLine($"data: {message}\n");
            }
            catch (IOException)
            {
                RemoveConnection(clientId);
            }
        }
    }
}

SSE实例的生命周期处理

  1. 连接管理

    • 使用线程安全集合存储活跃连接,避免多线程操作冲突。
    • 为每个连接分配唯一标识,方便定向推送或精准清理。
  2. 心跳机制
    定期发送心跳消息(注释或空数据),防止连接因长时间无数据被代理或浏览器断开:

// 定时心跳任务示例(使用BackgroundService)
public class SseHeartbeatService : BackgroundService
{
    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        while (!stoppingToken.IsCancellationRequested)
        {
            foreach (var kvp in _activeConnections)
            {
                try
                {
                    kvp.Value.WriteLine(": heartbeat\n");
                }
                catch (IOException)
                {
                    RemoveConnection(kvp.Key);
                }
            }
            await Task.Delay(30000, stoppingToken); // 每30秒发送一次心跳
        }
    }
}
  1. 断开清理

    • 监听请求的RequestAborted令牌,在客户端主动断开或请求超时后,从集合中移除连接并释放StreamWriter资源。
    • 在发送消息时捕获IOException(连接断开时抛出),自动清理无效连接。
  2. 资源释放
    确保StreamWriter在连接断开时被正确释放,避免内存泄漏,示例已包含在RemoveConnection方法中。

内容的提问来源于stack exchange,提问作者Francois Taljaard

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 16:35:32