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

泛型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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 01:10:32