C#并行线程调用异步函数出现任务未完成、文件乱码问题排查
代码问题分析
第一版代码任务未执行完成的核心错误
FileStream被提前释放:你将所有Task.Run任务的添加逻辑放在using包裹的fs作用域内,using块会在循环结束后立刻释放共享文件流,此时Task.WaitAll还没开始执行,所有访问已释放流的异步任务会直接抛出异常终止,无法走完完整流程。- 共享流并发操作无安全保障:单个
FileStream实例本身不是线程安全的,就算不提前释放,4个任务同时对同一个流做读、Seek、写操作,流内部的文件指针会互相抢占覆盖,既会出现读写位置错乱,也会触发流内部状态异常,导致任务意外退出。 - 闭包捕获变量风险:原代码中三个异步函数没有传入流、缓冲区、偏移等必要参数,如果闭包捕获的是循环外的共享变量,会出现多任务串用变量值的问题。
第二版代码写入乱码的核心原因
- 并发写无同步控制:虽然你给每个写操作单独创建了
FileStream,也开启了FileShare.ReadWrite允许文件共享访问,但操作系统文件层不保证非对齐并发写的原子性,多个流同时写入时会出现字节交错、部分内容被覆盖的问题,最终呈现为乱码。 - 偏移计算逻辑线程不安全:如果
getBlockOffset()是通过全局变量累加计算偏移,没有加锁或使用原子操作,多线程同时调用时会拿到重复的偏移值,两个任务往同一个文件位置写入数据,字节交错后会出现类似日文的乱码(本质是错位字节被解析为Unicode字符)。 - 读流被意外释放/修改:第二版代码中循环内的
using(fs = new FileStream(----))如果捕获的是外层声明的共享变量,会被循环反复新建、释放,异步任务执行时用到的读流可能已经被释放或移动了指针,读出来的缓冲区本身就是错误数据,写入后自然是乱码。
推荐替代实现方案
不要裸奔做多线程共享文件读写,优先采用「单线程拆分文件块+并行计算+单线程有序写入」的生产者消费者模式,既能保留并行处理的性能优势,又能从根源上避免文件并发操作的所有问题。
- 提前在单线程环境下计算好所有文件块的偏移、长度,不要把偏移计算放到多线程任务内。
- 用信号量
SemaphoreSlim控制并行任务的数量,每个任务读文件、计算时单独创建文件流,不要共享流实例。 - 所有计算完成后,单线程按偏移顺序写入结果文件,完全规避并发写的冲突问题。
- 异步代码中不要用
Task.WaitAll(),改用await Task.WhenAll(),避免阻塞线程池线程引发死锁。
参考实现代码:
int parallelCount = 4; int blockSize = 4 * 1024; // 对应你定义的单块长度 string inputFilePath = "你的输入文件路径"; string outputFilePath = "你的输出文件路径"; // 第一步:单线程提前拆分所有文件块,计算好每个块的偏移和读取长度 List<(long offset, int readLength)> fileBlocks = new(); using (var initReadStream = File.OpenRead(inputFilePath)) { long totalLength = initReadStream.Length; long currentOffset = 0; while (currentOffset < totalLength) { int currentReadLen = (int)Math.Min(blockSize, totalLength - currentOffset); fileBlocks.Add((currentOffset, currentReadLen)); currentOffset += currentReadLen; } } // 第二步:并行处理所有块的读、计算逻辑,用信号量控制并发度 SemaphoreSlim concurrencySemaphore = new SemaphoreSlim(parallelCount); ConcurrentBag<(long writeOffset, byte[] writeData)> writeQueue = new(); List<Task> processTasks = new(); foreach (var block in fileBlocks) { await concurrencySemaphore.WaitAsync(); processTasks.Add(Task.Run(async () => { try { byte[] readBuffer = new byte[block.readLength]; // 每个读操作单独创建流,不共享实例 using (var readStream = new FileStream(inputFilePath, FileMode.Open, FileAccess.Read, FileShare.Read)) { readStream.Seek(block.offset, SeekOrigin.Begin); await readStream.ReadExactlyAsync(readBuffer); } // 对应原Function2Async的计算逻辑 long computeResult = await DoComputeAsync(readBuffer); // 替换为你实际需要写入的字节内容 byte[] needWriteData = BitConverter.GetBytes(computeResult); writeQueue.Add((block.offset, needWriteData)); } finally { concurrencySemaphore.Release(); } })); } await Task.WhenAll(processTasks); // 第三步:所有计算完成后,单线程按顺序写入结果 using (var writeStream = new FileStream(outputFilePath, FileMode.Create, FileAccess.Write, FileShare.None)) { foreach (var writeItem in writeQueue.OrderBy(item => item.writeOffset)) { writeStream.Seek(writeItem.writeOffset, SeekOrigin.Begin); await writeStream.WriteAsync(writeItem.writeData); } }
如果确实需要边计算边写入,必须给写文件的逻辑加锁保证同一时间只有一个线程在操作写流,同时偏移计算必须用Interlocked.Add做原子操作,避免拿到重复偏移。
内容的提问来源于stack exchange,提问作者Logging device
相关产品推荐
相关产品推荐

