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

Parallel.Foreach循环中第三方SDK事件处理及数据库安全更新咨询

批量并行处理中的事件处理与数据库安全更新方案

问题描述

我正在对项目进行批量处理,采用如下并行处理方式。需要获取Task未返回时,由下方两个事件处理器(状态更新)传递的信息。想问:

  1. 这些事件能否在Parallel.ForEach循环中处理?
  2. 当前设置下是否会出现互相干扰的情况?
  3. 如何安全更新SQL数据库表中的状态?

现有代码

public async Task<Junk> blahblahblah(List<Item> unprocessedItems)
{          
    var tasks = new List<Task>();
    Parallel.ForEach(unprocessedItems,  item =>
    {
        var exporter = new Exporter();
        exporter.ExportEnded += Exporter_ExportEnded;
        exporter.StatisticsReceived += Exporter_StatisticsReceived;
        var exporterTask = exporter.ExportAsync(cameraConfig, PlaybackMode.Sequential, fileName, true);

        //do exporter task stuff
        tasks.Add(exporterTask);
    });

    await Task.WhenAll(tasks);
}
    
private void Exporter_StatisticsReceived(object sender, StatisticsEventArgs e)
{
    //do stuff with event information
    //update percent status in DB for each parallel row independently.
}
    
private void Exporter_ExportEnded(object sender, EndedEventArgs e)
{
    //do stuff with event information
    //log results in database.
}

问题解答

1. 事件能否在Parallel.ForEach中处理?

能处理,但必须解决线程安全问题。Parallel.ForEach会启用多个线程池线程执行循环体,Exporter的事件会在不同线程触发,多个线程同时进入事件处理器时,若没有同步机制,极易引发竞态条件。

2. 当前设置是否会出现干扰?

你的推测是对的,当前写法存在明显风险:

  • 事件处理器无法区分是哪个item触发的事件,可能导致数据库更新错误(比如把A项的进度更新到B项的记录里)。
  • 多个线程同时更新数据库时,会出现并发冲突,导致进度值覆盖、日志记录混乱等问题。
  • Parallel.ForEach不适合搭配异步方法(ExportAsync是IO密集型异步操作),会造成线程池资源浪费,且无法高效利用异步IO的优势。

3. 安全处理方案

(1)替换Parallel.ForEach为异步并行

针对IO密集型的导出操作,直接用foreach+Task.WhenAll更高效,避免不必要的线程占用:

public async Task<Junk> BlahBlahBlah(List<Item> unprocessedItems)
{          
    var tasks = new List<Task>();
    foreach (var item in unprocessedItems)
    {
        var exporter = new Exporter();
        // 捕获当前item,确保事件能关联到对应业务条目
        var currentItem = item;
        // 绑定带item参数的事件处理器
        exporter.ExportEnded += (sender, e) => Exporter_ExportEnded(sender, e, currentItem);
        exporter.StatisticsReceived += (sender, e) => Exporter_StatisticsReceived(sender, e, currentItem);
        
        var exporterTask = exporter.ExportAsync(cameraConfig, PlaybackMode.Sequential, fileName, true);
        tasks.Add(exporterTask);
    }

    await Task.WhenAll(tasks);
    return new Junk();
}

(2)事件处理器的线程安全与数据库并发控制

核心是确保每个事件对应唯一的业务条目,并通过数据库机制避免并发冲突:

进度更新(StatisticsReceived):用乐观并发控制

乐观并发适合高并发场景,通过检查当前数据库值是否未被修改来避免覆盖:

private async void Exporter_StatisticsReceived(object sender, StatisticsEventArgs e, Item item)
{
    using (var connection = new SqlConnection("你的连接字符串"))
    {
        await connection.OpenAsync();
        // 仅当数据库中当前进度与传入的旧进度一致时才更新
        var updateSql = @"UPDATE ExportTasks 
                          SET Progress = @newProgress 
                          WHERE Id = @itemId AND Progress = @oldProgress";
        var command = new SqlCommand(updateSql, connection);
        command.Parameters.AddWithValue("@newProgress", e.PercentComplete);
        command.Parameters.AddWithValue("@itemId", item.Id);
        command.Parameters.AddWithValue("@oldProgress", e.PreviousPercent); // 假设事件包含旧进度值
        
        var rowsAffected = await command.ExecuteNonQueryAsync();
        // 若影响行数为0,说明有其他线程已更新,可根据业务逻辑重试或忽略
    }
}
结果日志(ExportEnded):用唯一标识确保精准更新

每个导出任务对应唯一的item.Id,直接更新对应记录即可,若担心并发,可添加行级锁:

private async void Exporter_ExportEnded(object sender, EndedEventArgs e, Item item)
{
    using (var connection = new SqlConnection("你的连接字符串"))
    {
        await connection.OpenAsync();
        // 使用UPDLOCK行级锁,防止其他线程同时修改该记录
        var updateSql = @"UPDATE ExportTasks WITH (UPDLOCK)
                          SET Status = @status, EndTime = @endTime 
                          WHERE Id = @itemId";
        var command = new SqlCommand(updateSql, connection);
        command.Parameters.AddWithValue("@status", e.IsSuccess ? "Completed" : "Failed");
        command.Parameters.AddWithValue("@endTime", DateTime.UtcNow);
        command.Parameters.AddWithValue("@itemId", item.Id);
        
        await command.ExecuteNonQueryAsync();
    }
}

(3)其他注意事项

  • 若事件处理器需执行异步操作,应使用async void(仅适合事件处理器场景),并注意异常捕获(可在内部加try-catch避免未处理异常)。
  • 避免在事件处理器中访问非线程安全的共享变量,若必须使用,需用lock或线程安全集合(如ConcurrentDictionary)同步访问。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 01:40:27