.NET 7中Server-Sent Events连续发送时偶尔格式异常求助
问题分析与解决方案
核心问题:并发写入线程不安全导致内容交叉污染
你的假设完全正确,问题根源在于**Response.Body并非线程安全**,而你使用了async void类型的事件处理方法。当短时间内连续触发LobbyEvent或KeepAlive事件时,多个异步写入操作会同时抢占Response.Body资源,导致字节流交叉写入,最终出现Data字段异常、事件类型拼接错误(比如lobbyevent:lobbyEvent)等格式问题。线上环境响应速度更慢,并发冲突概率更高,因此问题表现更明显,本地难以复现也符合这一逻辑。
此外,async void本身存在异常处理困难、无法跟踪异步操作状态的缺陷,在事件处理场景中会进一步放大并发风险。
具体修复步骤
1. 为每个连接添加写入锁,保证原子性
在连接方法内部创建锁对象,确保同一时间只有一个写入操作能访问Response.Body,彻底避免交叉写入。
2. 替换async void事件处理为安全逻辑
将事件处理改为同步方法,内部通过锁保护异步写入操作,同时避免async void带来的潜在问题。
3. 优化SSE写入逻辑,减少IO操作次数
把单个SSE事件的所有内容拼接成完整字符串后一次性写入,减少多次WriteAsync调用带来的并发冲突概率。
修复后的完整代码
[Produces("text/event-stream")] [HttpGet("{code}/listen")] public async Task ListenForNotifications(string code, CancellationToken cancellationToken) { _logger.LogCritical($"CONNECTED"); var lobby = await _dbContext.Lobbies.FirstOrDefaultAsync(l => l.JoinCode == code); if (lobby == null) return; SetServerSentEventHeaders(); // 为当前连接创建写入锁,确保同一时间只有一个写入操作执行 var writeLock = new SemaphoreSlim(1, 1); // 封装SSE事件写入逻辑,统一处理锁和格式 async Task WriteSseEventAsync(string eventType, string? data, CancellationToken ct) { await writeLock.WaitAsync(ct); try { var sb = new StringBuilder(); sb.AppendLine($"event:{eventType}"); if (!string.IsNullOrEmpty(data)) { sb.AppendLine($"data:{data}"); } sb.AppendLine("retry:10000"); sb.AppendLine(); // SSE事件必须以空行结尾 await Response.WriteAsync(sb.ToString(), ct); await Response.Body.FlushAsync(ct); } finally { writeLock.Release(); } } // 连接成功事件 await WriteSseEventAsync("connected", null, cancellationToken); void OnNotification(object? sender, LobbyEventArgs eventArgs) { if (eventArgs.LobbyEvent.LobbyId != lobby.Id) return; try { var json = JsonSerializer.Serialize(eventArgs.LobbyEvent, _jsonSerializerOptions); // 异步执行写入,捕获异常避免崩溃 _ = WriteSseEventAsync("lobbyEvent", json, cancellationToken); _logger.LogCritical($"NOTIFICATION SENT: {eventArgs.LobbyEvent.Type}"); } catch (Exception ex) { _logger.LogCritical($"NOTIFICATION FAILED: {ex.Message}"); } } void KeepAlive(object? sender, EventArgs eventArgs) { try { _ = WriteSseEventAsync("keep-alive", "keep-alive", cancellationToken); } catch (Exception ex) { _logger.LogCritical($"KEEPALIVE FAILED: {ex.Message}"); } } _eventService.LobbyEvent += OnNotification; _eventService.KeepAlive += KeepAlive; try { while (!cancellationToken.IsCancellationRequested) await Task.Delay(1000, cancellationToken); } finally { _logger.LogCritical($"DISCONNECTED"); _eventService.LobbyEvent -= OnNotification; _eventService.KeepAlive -= KeepAlive; writeLock.Dispose(); } } private void SetServerSentEventHeaders() { Response.StatusCode = 200; Response.Headers.Add("Content-Type", "text/event-stream"); Response.Headers.Add("Cache-Control", "no-cache"); Response.Headers.Add("Connection", "keep-alive"); }
额外优化建议
- 引入消息队列:为每个连接维护一个消息队列,后台循环从队列读取消息并写入响应,能更精准控制写入节奏,避免事件处理逻辑阻塞。
- 检查取消状态:写入前先验证
CancellationToken是否已触发,避免执行无效IO操作。 - 使用成熟SSE库:.NET 7+可考虑使用官方SignalR或专门的SSE NuGet包,这类库已内置线程安全和格式校验逻辑,无需手动实现底层写入。
内容的提问来源于stack exchange,提问作者Santi
相关产品推荐
相关产品推荐

