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

如何在C#中无延迟插入WebSocket接收的数据

WebSocket高频数据入库延迟优化方案

问题背景

使用WebSocketSharp接收WebSocket数据,每秒触发约550次OnMessage事件,每秒需处理500条单行JSON数据,要求在1秒内完成数据库插入操作,但当前实现存在明显延迟。

数据示例:

{"seg":"BSE","currc":"0.012","code":"538610","ltt":"1683171716"}

当前代码实现:

WebSocket WebS = new WebSocket("SocketLinkHere");
WebS.OnOpen += (s1, e1) =>
{
    // 无逻辑
};
WebS.OnError += (s1, e1) =>
{
    if (!WebS.IsAlive)
    {
        WebS.Close();
        WebS.Connect();
    }
};
WebS.OnClose += (s1, e1) =>
{
    // 错误记录
};
WebS.OnMessage += (s1, e1) =>
{
    try
    {
        string result = "[" + e1.Data + "]";
        SQLEXE(result); // 调用插入方法
    }
    catch (Exception ex) { }
};

WebS.Connect();


private void SQLEXE(string result)
{
    try
    {
        DataTable dt = JsonConvert.DeserializeObject<DataTable>(result);
        SqlParameter[] para = {
            new SqlParameter("@ltt", dt.Rows[0]["ltt"]),
            new SqlParameter("@code", dt.Rows[0]["code"]),
            new SqlParameter("@currc", dt.Rows[0]["currc"]),
            new SqlParameter("@seg", dt.Rows[0]["seg"])
        };
        cmd.Parameters.Clear();
        cmd.Parameters.AddRange(para);
        cmd.ExecuteNonQuery();
    }
    catch (Exception error) { }
}

当前代码的核心性能瓶颈

  • 单条插入+同步阻塞:每次OnMessage触发都同步执行一次数据库插入,500次/秒的频率会导致大量数据库连接开销和IO阻塞,完全无法满足毫秒级要求。
  • JSON反序列化冗余:将单条JSON包装成数组再反序列化为DataTable,额外增加了序列化开销,完全没必要。
  • 资源复用不规范:cmd对象的复用方式存在风险,且未利用数据库连接池优化连接的创建与销毁。

优化方案

1. 批量插入+异步缓冲

将高频接收的数据先缓冲到内存队列,定时(比如每50ms)或累积到指定数量(比如100条)时执行批量插入,大幅减少数据库交互次数。

2. 优化JSON反序列化

直接将单条JSON反序列化为实体类,避免DataTable的冗余开销:

public class MarketData
{
    public string seg { get; set; }
    public string currc { get; set; }
    public string code { get; set; }
    public string ltt { get; set; }
}

3. 使用SqlBulkCopy高效批量入库

利用SqlServer的SqlBulkCopy组件实现高速批量插入,比单条Insert语句效率提升数倍。

4. 线程安全队列+后台处理

使用ConcurrentQueue存储待插入数据,后台开启独立线程处理批量入库,避免阻塞WebSocket的接收线程。

完整优化代码示例

using System.Collections.Concurrent;
using System.Data.SqlClient;
using Newtonsoft.Json;
using WebSocketSharp;

public class MarketDataProcessor
{
    private WebSocket _webSocket;
    private readonly ConcurrentQueue<MarketData> _dataQueue = new ConcurrentQueue<MarketData>();
    private readonly Timer _batchInsertTimer;
    private const int BatchSize = 100; // 每累积100条执行一次批量插入
    private const int IntervalMs = 50; // 每50ms检查一次队列

    public MarketDataProcessor(string socketUrl)
    {
        _webSocket = new WebSocket(socketUrl);
        _webSocket.OnMessage += WebSocket_OnMessage;
        _webSocket.OnError += WebSocket_OnError;
        _webSocket.OnClose += WebSocket_OnClose;

        // 初始化批量插入定时器
        _batchInsertTimer = new Timer(ProcessBatchInsert, null, 0, IntervalMs);
    }

    public void Start()
    {
        _webSocket.Connect();
    }

    private void WebSocket_OnMessage(object sender, MessageEventArgs e)
    {
        try
        {
            // 直接反序列化为实体类,跳过冗余包装和DataTable转换
            var data = JsonConvert.DeserializeObject<MarketData>(e.Data);
            if (data != null)
            {
                _dataQueue.Enqueue(data);
            }
        }
        catch (Exception ex)
        {
            // 记录日志,禁止吞掉异常
        }
    }

    private void WebSocket_OnError(object sender, ErrorEventArgs e)
    {
        if (!_webSocket.IsAlive)
        {
            _webSocket.Close();
            _webSocket.Connect();
        }
    }

    private void WebSocket_OnClose(object sender, CloseEventArgs e)
    {
        // 记录连接关闭日志
    }

    private void ProcessBatchInsert(object state)
    {
        if (_dataQueue.Count == 0) return;

        List<MarketData> batchData = new List<MarketData>();
        // 从队列中取出批量数据
        while (batchData.Count < BatchSize && _dataQueue.TryDequeue(out var data))
        {
            batchData.Add(data);
        }

        if (batchData.Count == 0) return;

        try
        {
            using (var conn = new SqlConnection("YourDatabaseConnectionString"))
            {
                conn.Open();
                using (var bulkCopy = new SqlBulkCopy(conn))
                {
                    bulkCopy.DestinationTableName = "YourTableName"; // 替换为你的目标表名
                    // 映射实体类字段与数据库表字段
                    bulkCopy.ColumnMappings.Add("seg", "seg");
                    bulkCopy.ColumnMappings.Add("currc", "currc");
                    bulkCopy.ColumnMappings.Add("code", "code");
                    bulkCopy.ColumnMappings.Add("ltt", "ltt");

                    // 将List转为DataTable供SqlBulkCopy使用
                    var dt = new DataTable();
                    dt.Columns.Add("seg", typeof(string));
                    dt.Columns.Add("currc", typeof(string));
                    dt.Columns.Add("code", typeof(string));
                    dt.Columns.Add("ltt", typeof(string));

                    foreach (var item in batchData)
                    {
                        dt.Rows.Add(item.seg, item.currc, item.code, item.ltt);
                    }

                    bulkCopy.WriteToServer(dt);
                }
            }
        }
        catch (Exception ex)
        {
            // 记录批量插入异常,可考虑将失败数据重新入队重试
        }
    }
}

public class MarketData
{
    public string seg { get; set; }
    public string currc { get; set; }
    public string code { get; set; }
    public string ltt { get; set; }
}

额外优化建议

  • 配置数据库连接池:在连接字符串中调整Max Pool Size参数,确保有足够的连接处理批量操作。
  • 禁止吞掉异常:添加日志框架记录异常,便于排查问题。
  • 字段类型匹配:确保实体类字段类型与数据库表字段类型完全一致,避免不必要的类型转换开销。
  • 队列容量控制:若数据产生速度远大于入库速度,需添加队列溢出处理逻辑,防止内存溢出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 17:47:06