.NET 6中如何正确实现List<List<Foo>>的Parallel.ForEachAsync
正确实现基于List<List>的Parallel.ForEachAsync并行数据库写入
现有代码的核心问题
- 线程不安全的计数操作:
rows += updatedRows是非原子操作,多个并行任务同时修改会导致最终计数不准确。 - SQL构建错误:代码里有变量名错误(
sp.Append应为sb.Append),而且SQL拼接格式不对——每个条目没闭合括号、多个条目之间没有逗号分隔,这会引发SQL执行异常或数据插入错误。 - 未控制并行度:默认并行度可能过高,导致数据库连接池耗尽或并发写入冲突,进而出现数据重复或异常。
- REPLACE INTO逻辑隐患:如果没基于主键/唯一键正确执行替换,并行写入时的竞态条件可能导致重复条目插入。
修正后的实现代码
using System.Threading; var totalRows = 0L; var parallelOptions = new ParallelOptions { // 根据数据库连接池大小和服务器性能调整,避免过度并发 MaxDegreeOfParallelism = Math.Min(Environment.ProcessorCount * 2, 10) }; await Parallel.ForEachAsync(listOfLists, parallelOptions, async (list, token) => { var sb = new StringBuilder("REPLACE INTO `table` VALUES "); var isFirst = true; foreach (var item in list) { if (!isFirst) { sb.Append(","); } // 注意:根据实际字段数量、类型调整格式,这里假设item.bar是单个字段,同时做SQL转义避免注入 sb.Append($"({EscapeSqlValue(item.bar)})"); isFirst = false; } // 避免空列表生成无效SQL if (sb.Length > "REPLACE INTO `table` VALUES ".Length) { var updatedRows = await _client.ExecuteSql(sb.ToString(), token); // 用Interlocked.Add保证线程安全的计数累加 Interlocked.Add(ref totalRows, updatedRows); } }); // 辅助方法:转义SQL值,避免注入和格式错误 string EscapeSqlValue(object value) { if (value is string str) { return $"'{str.Replace("'", "''")}'"; } return value?.ToString() ?? "NULL"; }
关键修正点说明
- 线程安全计数:用
Interlocked.Add替代直接的+=操作,保证多任务修改计数时的原子性,避免计数丢失或错误。 - 规范SQL拼接:
- 修复变量名错误,确保每个并行任务使用独立的
StringBuilder实例。 - 给每个条目添加闭合括号,多个条目间用逗号分隔,避免SQL语法错误。
- 增加空列表判断,防止生成无效SQL语句。
- 加入SQL值转义逻辑,规避注入风险和数据格式问题。
- 修复变量名错误,确保每个并行任务使用独立的
- 控制并行度:通过
ParallelOptions.MaxDegreeOfParallelism限制并发任务数,平衡写入性能和数据库压力,避免连接池耗尽或并发冲突。 - 传递取消令牌:把
token传给ExecuteSql方法,支持任务取消,提升系统响应性。
重复条目问题的额外排查方向
如果修正后仍有重复条目,需要检查:
- REPLACE INTO的唯一键设置:确认目标表已正确定义主键或唯一约束,REPLACE INTO只会替换与唯一键匹配的条目,否则会插入新行。
- 数据拆分逻辑:检查
listOfLists的拆分过程,确保没有把同一个Foo条目重复分配到多个子列表,导致并行写入时重复插入。 - 数据库隔离级别:若隔离级别过低,可能引发并发写入时的幻读或重复插入,可适当调整为READ COMMITTED等合适的级别。
内容的提问来源于stack exchange,提问作者lbueno
相关产品推荐
相关产品推荐

