如何在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
相关产品推荐
相关产品推荐

