并行循环中数据插入实现及BufferedStream、并行数设置疑问
嘿,针对你处理大文本文件并行插入数据库的需求,我来帮你梳理清楚几个核心问题,并且修正你的代码:
大文本文件并行处理与数据库批量插入解决方案
1. 关于BufferedStream的使用
其实StreamReader默认已经内置了缓冲区(默认1024字节,可通过构造函数调整大小),直接用它读取大文件已经足够高效。如果你的场景需要更精细地控制缓冲区大小(比如针对500MB级文件优化IO性能),再用BufferedStream包装FileStream才是有意义的,比如设置64KB或128KB的缓冲区:
using (var fs = File.Open(filePath, FileMode.Open, FileAccess.Read, FileShare.Read)) using (var bs = new BufferedStream(fs, 65536)) // 64KB缓冲区 using (var sr = new StreamReader(bs)) { // 读取逻辑 }
如果没有特殊需求,直接用new StreamReader(filePath)就够,底层已经做了缓冲优化,不用画蛇添足。
2. 并行操作数的设置
你的场景是IO密集型任务(读文件+数据库写入),不是CPU密集型,完全不能依赖Parallel.For的默认设置(默认会根据CPU核心数开多线程)——不然很容易导致磁盘IO竞争、数据库连接池耗尽。
- 你现在只有2个文件,Parallel.For默认会并行处理这2个(文件数少于CPU核心数),无需额外设置;
- 如果后续文件数量增加,一定要通过
ParallelOptions限制并行度,比如设置为CPU核心数或者更小的值(比如4),避免资源过载:
var parallelOptions = new ParallelOptions { MaxDegreeOfParallelism = Math.Min(Environment.ProcessorCount, 4) // 最多4个并行任务 }; Parallel.For(0, SourceFiles.Length, parallelOptions, x => { // 单个文件处理逻辑 });
另外,数据库连接池默认最大连接数是100,所以并行数建议远小于这个值(比如10以内),防止连接池耗尽。
3. 正确实现数据插入(结合SqlBulkCopy)
你的初始代码有几个明显问题:
ProcessAllLinesInInputFile里的Parallel.For(0, SourceFiles, x =>是错误的(SourceFiles是字符串,不是长度);DataTable不是线程安全的,不能在并行循环中直接修改;- 没有实现分批次提交逻辑,一次性加载整个大文件到内存会直接爆内存。
核心优化思路
不要并行处理单个文件内的行!磁盘IO是顺序型操作,并行读取同一文件的不同行会导致频繁磁盘寻道,反而变慢。正确的做法是:单线程顺序读取行,分批次填充DataTable,达到指定批次大小后批量写入数据库,再清空DataTable继续处理。
修正后的完整代码
using System; using System.Data; using System.Data.SqlClient; using System.IO; using System.Threading.Tasks; class Program { static void Main(string[] args) { ProcessFileTaskItem( new string[] { @".\Insert.txt", @".\Insert1.txt" }, @"Data Source=(localdb)\MSSQLLocalDB;Initial Catalog=test;Integrated Security=True;Connect Timeout=30;Encrypt=False;TrustServerCertificate=False;ApplicationIntent=ReadWrite;MultiSubnetFailover=False", "test" ); } /// <summary> /// 并行处理多个文件,批量插入数据库 /// </summary> public static void ProcessFileTaskItem(string[] sourceFiles, string dbConnectionString, string destinationTable) { if (sourceFiles == null || sourceFiles.Length == 0) return; // 针对2个文件,限制并行度为2即可,避免不必要的线程切换 var parallelOptions = new ParallelOptions { MaxDegreeOfParallelism = 2 }; Parallel.For(0, sourceFiles.Length, parallelOptions, fileIndex => { var filePath = sourceFiles[fileIndex]; if (!File.Exists(filePath)) { Console.WriteLine($"文件不存在:{filePath}"); return; } // 每个文件用独立的数据库连接,依赖连接池自动管理 using (var connection = new SqlConnection(dbConnectionString)) { connection.Open(); using (var bulkCopy = new SqlBulkCopy(connection, SqlBulkCopyOptions.TableLock, null)) { bulkCopy.DestinationTableName = destinationTable; bulkCopy.BulkCopyTimeout = 28800; // 8小时超时 bulkCopy.BatchSize = 10000; // 每10000行提交一次,可根据内存调整 // 字段映射(如果数据库列名和解析字段名完全一致,可省略这一步) bulkCopy.ColumnMappings.Add("Name", "Name"); bulkCopy.ColumnMappings.Add("Comment", "Comment"); bulkCopy.ColumnMappings.Add("Address", "Address"); bulkCopy.ColumnMappings.Add("Phone", "Phone"); bulkCopy.ColumnMappings.Add("IsActive", "IsActive"); // 分批次处理文件内容 ProcessFileLinesInBatches(filePath, bulkCopy); } } }); Array.Clear(sourceFiles, 0, sourceFiles.Length); Console.WriteLine("所有文件处理完成"); } /// <summary> /// 分批次读取文件行,批量插入数据库 /// </summary> private static void ProcessFileLinesInBatches(string filePath, SqlBulkCopy bulkCopy) { // 创建与数据库表结构匹配的DataTable var batchTable = new DataTable(); batchTable.Columns.Add("Name", typeof(string)); batchTable.Columns.Add("Comment", typeof(string)); batchTable.Columns.Add("Address", typeof(string)); batchTable.Columns.Add("Phone", typeof(string)); batchTable.Columns.Add("IsActive", typeof(int)); const int batchSize = 10000; int currentRowCount = 0; // 直接用StreamReader读取,默认缓冲足够高效 using (var sr = new StreamReader(filePath)) { string line; while ((line = sr.ReadLine()) != null) { // 解析每行数据:处理带引号的字段,分割后清理空格和引号 var fields = line.Split(new[] { ',' }, StringSplitOptions.RemoveEmptyEntries); if (fields.Length != 5) { Console.WriteLine($"无效行,跳过:{line}"); continue; } var name = fields[0].Trim('"', ' '); var comment = fields[1].Trim('"', ' '); var address = fields[2].Trim('"', ' '); var phone = fields[3].Trim('"', ' '); if (!int.TryParse(fields[4].Trim('"', ' '), out int isActive)) { Console.WriteLine($"IsActive字段解析失败,跳过:{line}"); continue; } // 添加到批次表 batchTable.Rows.Add(name, comment, address, phone, isActive); currentRowCount++; // 达到批次大小,提交数据 if (currentRowCount >= batchSize) { bulkCopy.WriteToServer(batchTable); batchTable.Clear(); currentRowCount = 0; } } } // 提交剩余的未达批次的行 if (currentRowCount > 0) { bulkCopy.WriteToServer(batchTable); batchTable.Clear(); } } }
关键注意事项
- 线程安全:每个文件的处理完全独立,DataTable仅在单个线程内使用,避免了线程安全问题;
- 内存控制:通过
BatchSize分批次提交,不会一次性加载整个大文件到内存,避免内存溢出; - 连接管理:
using语句会自动释放SqlConnection,依赖ADO.NET连接池自动复用连接,不用手动调用Close(); - 性能优化:使用
SqlBulkCopyOptions.TableLock可以提升批量插入速度(前提是表没有其他并发写入操作); - 错误处理:代码中添加了基本的行解析错误处理,你可以根据需求扩展日志或重试逻辑。
内容的提问来源于stack exchange,提问作者 cickness
相关产品推荐
相关产品推荐

