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

C# .NET 8 SSE服务器有响应但浏览器不显示消息,求排查

问题分析

你的SSE服务器数据已发送到客户端,但JavaScript的EventSource无响应,同时curl测试出现数据缓存后批量输出的问题,核心原因是服务器端未正确异步刷新输出流,导致数据被操作系统或.NET流缓冲区缓存,无法实时推送给客户端。EventSource需要接收到符合规范的完整消息块(以\n\n结尾)才会触发onmessage,缓冲导致消息无法及时到达。

修复后的C#服务器代码
using System;
using System.Collections.Generic;
using System.Net;
using System.Text;
using System.Threading.Tasks;
using System.Timers;

class Program
{
    private static readonly List<HttpListenerContext> clients = new List<HttpListenerContext>();
    private static readonly object lockObject = new object();
    private static System.Timers.Timer? messageTimer;

    static async Task Main(string[] args)
    {
        string url = "http://+:5000/";
        HttpListener listener = new HttpListener();
        listener.Prefixes.Add(url);
        listener.Start();
        Console.WriteLine($"Server runs on {url}");

        // 定时发送消息
        messageTimer = new System.Timers.Timer(1000);
        messageTimer.Elapsed += SendPeriodicMessages;
        messageTimer.AutoReset = true;
        messageTimer.Enabled = true;

        while (true)
        {
            HttpListenerContext context = await listener.GetContextAsync();
            await HandleClient(context); // 异步处理客户端连接
        }
    }

    // 改为异步方法,支持await刷新操作
    private static async Task HandleClient(HttpListenerContext context)
    {
        HttpListenerRequest request = context.Request;
        HttpListenerResponse response = context.Response;

        // CORS配置
        response.Headers.Add("Access-Control-Allow-Origin", "*");
        response.Headers.Add("Access-Control-Allow-Methods", "GET, OPTIONS");
        response.Headers.Add("Access-Control-Allow-Headers", "Content-Type");

        if (request.HttpMethod == "OPTIONS")
        {
            response.StatusCode = (int)HttpStatusCode.OK;
            response.OutputStream.Close();
            return;
        }

        // SSE响应头配置
        response.ContentType = "text/event-stream";
        response.Headers.Add("Cache-Control", "no-cache");
        response.Headers.Add("Connection", "keep-alive");
        response.SendChunked = true; // 明确启用分块传输

        // 发送初始心跳消息,激活客户端连接(符合SSE规范的注释行,不会触发onmessage)
        var initBuffer = Encoding.UTF8.GetBytes(": ping\n\n");
        await response.OutputStream.WriteAsync(initBuffer, 0, initBuffer.Length);
        await response.OutputStream.FlushAsync();

        AddClient(context);
    }

    private static void AddClient(HttpListenerContext context)
    {
        lock (lockObject)
        {
            clients.Add(context);
            var clientIp = context.Request.RemoteEndPoint.Address.ToString();
            var clientPort = context.Request.RemoteEndPoint.Port;
            Console.WriteLine($"Client connected: IP = {clientIp}, Port = {clientPort}");
        }
    }

    static int TimerCounter = 0;
    private static void SendPeriodicMessages(object? sender, ElapsedEventArgs e)
    {
        Console.WriteLine("TimerTick " + TimerCounter);
        _ = SendMessagesToClients($"data: {TimerCounter++}\n\n"); // 火并忘却,不阻塞定时器
    }

    private static async Task SendMessagesToClients(string message)
    {
        byte[] buffer = Encoding.UTF8.GetBytes(message);
        List<Task> sendTasks = new List<Task>();
        List<HttpListenerContext> removeList = new();

        lock (lockObject)
        {
            Console.WriteLine("Number of Clients: " + clients.Count);
            foreach (var client in clients)
            {
                sendTasks.Add(SendToSingleClient(client, buffer, removeList));
            }
        }

        await Task.WhenAll(sendTasks);

        lock (lockObject)
        {
            foreach (var client in removeList)
            {
                clients.Remove(client);
            }
        }
    }

    // 单独抽离客户端发送逻辑,异步执行
    private static async Task SendToSingleClient(HttpListenerContext client, byte[] buffer, List<HttpListenerContext> removeList)
    {
        try
        {
            var clientIp = client.Request.RemoteEndPoint.Address.ToString();
            var clientPort = client.Request.RemoteEndPoint.Port;
            Console.WriteLine($"Sending Data ({buffer.Length}) to {clientIp}:{clientPort}");

            await client.Response.OutputStream.WriteAsync(buffer, 0, buffer.Length);
            await client.Response.OutputStream.FlushAsync(); // 异步刷新,确保数据实时发送
        }
        catch (Exception ex)
        {
            Console.WriteLine($"Error - Cant send data to Client: {ex.Message}");
            lock (lockObject) // 确保线程安全添加到移除列表
            {
                removeList.Add(client);
            }
        }
    }
}
关键改动说明
  1. 异步化处理:将HandleClient改为异步方法,调用WriteAsync和FlushAsync并等待完成,确保数据被立即推送到网络,而非留在缓冲区。
  2. 启用分块传输:明确设置response.SendChunked = true,配合SSE的流式传输需求,避免内容长度缓存。
  3. 初始心跳消息:发送: ping\n\n(SSE规范中的注释消息),让客户端确认连接建立,避免EventSource因长时间无数据而断开。
  4. 优化发送逻辑:抽离单客户端发送逻辑,减少线程池嵌套开销,同时在异常处理中确保移除列表的线程安全。
  5. 简化CORS配置:OPTIONS请求中无需重复添加CORS头,已在开头统一设置。
客户端代码(无需修改)
<!DOCTYPE html>
<html lang="en">
<head>
    <meta charset="UTF-8">
    <meta http-equiv="X-UA-Compatible" content="IE=edge">
    <meta name="viewport" content="width=device-width, initial-scale=1.0">
    <title>TEST</title>
</head>
<body>
    <h1>
        Server Send Event
    </h1>
    <div id="messages"></div>
    <script>
        const eventSource = new EventSource('http://192.168.56.245:5000/');

        eventSource.onmessage = function(event) {
            const messagesDiv = document.getElementById('messages');
            messagesDiv.innerHTML += `<p>${event.data}</p>`;
            console.log(event.data);
        };

        eventSource.onerror = function(event) {
            console.error("Error receiving messages from SSE:", event);
            eventSource.close();
        };
    </script>
</body>
</html>

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 23:05:54