能否在响应过滤器中使用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事件的步骤
初始化与连接
服务器端初始化
- 注册响应过滤器(或中间件),识别SSE请求并设置响应头。
- 维护一个线程安全集合(如
ConcurrentDictionary<string, StreamWriter>)存储活跃连接,关联客户端唯一标识(如SessionId、自定义ID)。 - 监听请求取消/断开事件,及时清理无效连接。
客户端连接
使用浏览器原生EventSourceAPI建立连接:
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实例的生命周期处理
连接管理
- 使用线程安全集合存储活跃连接,避免多线程操作冲突。
- 为每个连接分配唯一标识,方便定向推送或精准清理。
心跳机制
定期发送心跳消息(注释或空数据),防止连接因长时间无数据被代理或浏览器断开:
// 定时心跳任务示例(使用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秒发送一次心跳 } } }
断开清理
- 监听请求的
RequestAborted令牌,在客户端主动断开或请求超时后,从集合中移除连接并释放StreamWriter资源。 - 在发送消息时捕获
IOException(连接断开时抛出),自动清理无效连接。
- 监听请求的
资源释放
确保StreamWriter在连接断开时被正确释放,避免内存泄漏,示例已包含在RemoveConnection方法中。
内容的提问来源于stack exchange,提问作者Francois Taljaard
相关产品推荐
相关产品推荐

