泛型List并发访问问题:数据写入时清理列表
WebSocket流式数据无锁批量入库解决方案
1. 双缓冲列表切换方案
维护两个独立的List<T>实例,分别承担实时数据接收和待入库数据存储的角色,通过原子交换实现无锁的读写分离:
- 一个列表作为「活跃列表」,专门接收WebSocket推送的实时数据,添加数据时仅用粒度极小的锁保证单线程写入(仅
Add操作加锁,耗时可忽略)。 - 当需要入库时,用原子操作(如
Interlocked.Exchange)将活跃列表与空的待入库列表交换,此时新的流式数据会自动写入新的活跃列表,旧列表则可以安全枚举写入数据库,完全不阻塞数据接收。
示例代码(C#):
private List<T> _activeList = new List<T>(); private List<T> _persistList = new List<T>(); // WebSocket数据回调 public void OnWebSocketDataReceived(T data) { lock (_activeList) { _activeList.Add(data); } } // 定时/阈值触发的入库任务 public async Task PersistBatchAsync() { // 原子交换两个列表,拿到待入库的数据集 var batchData = Interlocked.Exchange(ref _activeList, _persistList); // 重置待入库列表,准备下次交换 _persistList = new List<T>(); // 安全写入数据库,此时batchData不会再被修改 foreach (var item in batchData) { await _dbContext.AddAsync(item); } await _dbContext.SaveChangesAsync(); // 释放内存 batchData.Clear(); }
2. 用线程安全队列替代普通List
直接使用ConcurrentQueue<T>替代List<T>,它原生支持并发入队/出队操作,无需手动加锁:
- WebSocket推送的数据直接调用
Enqueue写入队列,全程无阻塞。 - 入库任务批量取出队列中所有元素(通过
TryDequeue循环),再一次性写入数据库,避免频繁操作队列。
示例代码(C#):
private ConcurrentQueue<T> _dataQueue = new ConcurrentQueue<T>(); // WebSocket数据接收 public void OnWebSocketDataReceived(T data) { _dataQueue.Enqueue(data); } // 入库任务 public async Task PersistQueueDataAsync() { var batchItems = new List<T>(); // 一次性取出所有队列元素 while (_dataQueue.TryDequeue(out var item)) { batchItems.Add(item); } if (batchItems.Count == 0) return; await _dbContext.AddRangeAsync(batchItems); await _dbContext.SaveChangesAsync(); }
3. 快照复制+分段批量处理
当列表数据达到预设阈值时,快速复制一份当前列表的快照,清空原列表继续接收数据,再用快照进行入库操作:
- 添加数据时加短锁,仅在复制快照和清空列表的瞬间占用锁,对实时写入的影响微乎其微。
- 快照复制完成后,原列表立即恢复接收数据,入库操作在后台异步执行,完全不阻塞流式数据推送。
示例代码(C#):
private List<T> _activeList = new List<T>(); private readonly object _lockObj = new object(); private const int BatchThreshold = 1000; // 批量入库阈值 public void OnWebSocketDataReceived(T data) { lock (_lockObj) { _activeList.Add(data); // 达到阈值触发异步入库 if (_activeList.Count >= BatchThreshold) { _ = PersistSnapshotAsync(); } } } private async Task PersistSnapshotAsync() { List<T> snapshot; lock (_lockObj) { snapshot = _activeList.ToList(); _activeList.Clear(); } await _dbContext.AddRangeAsync(snapshot); await _dbContext.SaveChangesAsync(); }
内容的提问来源于stack exchange,提问作者Aditya Bokade
相关产品推荐
相关产品推荐

