多线程文件写入数组越界异常问题及优化方案咨询
问题
我开发了一款处理**大型数据集(LARGE datasets)**的应用,会对数据集执行各类转换后输出。该处理过程对时间要求极高,因此已投入大量精力进行优化。
设计思路为:每次读取一批记录,在不同线程中处理每条记录,并将结果写入文件。为尽可能避免内存写入保护异常或性能瓶颈,结果并非写入单个文件,而是写入多个临时文件,最终合并为目标输出文件。
为实现此方案,我创建了包含10个fileUtils的数组,线程初始化时会分配其中一个。通过threadedOutputIterator在每次localInit时递增,当计数达到10时重置为0,以此决定为每个线程的记录处理对象分配哪个fileUtils,每个util类负责收集并写入一个临时输出文件。
需注意,每个FileUtils对象会先在成员变量outputBuildString中收集约100条记录后再写入文件,因此这些对象独立于线程流程存在(线程内对象生命周期有限)。
该方案旨在将输出数据的收集、存储及写入职责尽可能均匀地分配给多个fileUtils对象,以实现比单文件写入更高的每秒写入量。
但当前方案出现了**数组越界(Array Out Of Bounds)**异常:尽管已有代码在threadedOutputIterator达到上限时将其重置,该值仍会超出上限(默认threadCount=10),导致访问threadedFileUtils数组(仅含10个对象)时触发异常。
代码如下(默认threadCount = 10):
private void ProcessRecords() { try { Parallel.ForEach(clientInputRecordList, new ParallelOptions { MaxDegreeOfParallelism = threadCount }, LocalInit, ThreadMain, LocalFinally); } catch (Exception e) { Console.WriteLine("The following error occured: " + e); } } private SplitLineParseObject LocalInit() { if (threadedOutputIterator >= threadCount) { threadedOutputIterator = 0; } // 仍会超出10,此处触发异常,因为threadedFileUtils数组只有10个对象 SplitLineParseObject splitLineParseUtil = new SplitLineParseObject(parmUtils, ref recCount, ref threadedFileUtils[threadedOutputIterator], ref recordsPassedToFileUtils); if (threadedOutputIterator < threadCount) { threadedOutputIterator++; } return splitLineParseUtil; } private SplitLineParseObject ThreadMain(ClientInputRecord record, ParallelLoopState state, SplitLineParseObject threadLocalObject) { threadLocalObject.clientInputRecord = record; threadLocalObject.ProcessRecord(); recordsPassedToObject++; return threadLocalObject; } private void LocalFinally(SplitLineParseObject obj) { obj = null; }
问题根源在于多线程并发访问threadedOutputIterator:多个线程可能同时执行判断和递增操作,导致该值突破上限,触发数组越界。需要优化方案,既要避免数组越界异常,又要保留多fileUtils带来的读写效率优势。
优化方案
1. 原子操作保证索引的线程安全
threadedOutputIterator的非原子读写是竞态条件的核心,用Interlocked类的原子方法替代手动判断和递增,再通过取模运算确保索引合法:
修改LocalInit方法如下:
private SplitLineParseObject LocalInit() { // 原子递增并获取当前值,减1是因为Increment会先加1再返回 int currentIndex = Interlocked.Increment(ref threadedOutputIterator) - 1; // 取模得到0到threadCount-1的合法索引 currentIndex = currentIndex % threadCount; // 处理int溢出导致的负数情况(实际场景中极少发生) if (currentIndex < 0) currentIndex += threadCount; SplitLineParseObject splitLineParseUtil = new SplitLineParseObject(parmUtils, ref recCount, ref threadedFileUtils[currentIndex], ref recordsPassedToFileUtils); return splitLineParseUtil; }
核心逻辑:Interlocked.Increment确保递增操作不会被线程打断,取模运算直接限制索引范围,从根源上杜绝数组越界。
2. 线程绑定FileUtils(可选,进一步优化并发效率)
如果希望每个线程固定使用同一个FileUtils,减少跨线程访问同一个FileUtils的冲突,可以用ThreadLocal<T>存储线程专属索引:
// 类级别声明ThreadLocal变量,线程首次初始化时分配唯一索引 private ThreadLocal<int> _threadFileIndex = new ThreadLocal<int>(() => { int index = Interlocked.Increment(ref threadedOutputIterator) % threadCount; return index < 0 ? index + threadCount : index; }); private SplitLineParseObject LocalInit() { int currentIndex = _threadFileIndex.Value; SplitLineParseObject splitLineParseUtil = new SplitLineParseObject(parmUtils, ref recCount, ref threadedFileUtils[currentIndex], ref recordsPassedToFileUtils); return splitLineParseUtil; }
优势:每个线程仅在首次初始化时获取一次索引,后续复用同一个FileUtils,降低FileUtils内部的并发写入冲突,进一步提升写入效率。
3. 补充:FileUtils自身的线程安全校验
FileUtils内部的outputBuildString如果是普通StringBuilder,多线程写入会导致数据错乱,需补充线程安全保障:
- 在FileUtils的记录收集、文件写入方法中添加
lock语句 - 或者替换为线程安全的缓存结构(如自行实现带锁的字符串构建逻辑)
内容的提问来源于stack exchange,提问作者Glenncito

