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

.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 03:45:41