C# .NET中锁定数组内指定分区解决并行写入数据损坏问题
解决ADLS Gen1 ConcurrentAppend并行写入的数据损坏问题
你遇到的核心问题是共享变量的竞态条件和无保护的缓冲区访问,多个线程同时修改counter、读写buffer时出现了位置混乱和数据重叠,最终导致写入ADLS的文件损坏。下面是针对性的修复方案:
问题根源拆解
你的代码里有几个关键的线程不安全点:
- 共享的
counter变量被多个线程同时读写,counter += ...这类操作是非原子的,会导致位置计算完全错乱。 - 对
buffer的读写、ConcurrentAppend调用、Array.Clear操作没有同步机制,即便操作不同分区,也可能因为counter错误写入到错误区域。 - 文件流没有用安全的方式释放,容易引发资源泄漏。
修复方案
我们要做到分区独立控制+细粒度锁,既保证并行效率,又避免临界区冲突:
完整修复代码
// 定义分区数量 var partitionCount = 5; // 为每个分区创建独立的锁对象(只锁当前操作的分区,不影响其他线程) var partitionLocks = new object[partitionCount]; for (int i = 0; i < partitionCount; i++) { partitionLocks[i] = new object(); } // 预先收集每个文件的路径和大小,避免并行时重复IO var fileDetails = new (string Path, long Size)[partitionCount]; for (int i = 0; i < partitionCount; i++) { var path = @"C:\Users\t-chkum\Desktop\InputFiles\1MB\" + (i+1) + ".txt"; fileDetails[i] = (path, new FileInfo(path).Length); } // 预先计算每个分区在buffer中的固定起始位置,彻底抛弃共享counter var partitionOffsets = new int[partitionCount]; int currentOffset = 0; for (int i = 0; i < partitionCount; i++) { partitionOffsets[i] = currentOffset; currentOffset += (int)fileDetails[i].Size + 1; // 和你原来的间隔逻辑保持一致 } Parallel.For(0, partitionCount, i => { var (filePath, fileSize) = fileDetails[i]; var currentPartitionOffset = partitionOffsets[i]; var currentLock = partitionLocks[i]; // 使用using自动释放文件流,避免资源泄漏 using (var stream = File.OpenRead(filePath)) { // 只锁定当前分区对应的锁,其他分区可以正常并行操作 lock (currentLock) { // 读取数据到当前分区的固定buffer位置 stream.Read(buffer, currentPartitionOffset, (int)fileSize); // 调用ConcurrentAppend,使用预先计算好的偏移量 client.ConcurrentAppend(fileName, true, buffer, currentPartitionOffset, (int)fileSize); // 清理当前分区的buffer(如果后续要复用buffer的话) Array.Clear(buffer, currentPartitionOffset, (int)fileSize); } } // 如果需要循环复用buffer,这里可以针对单个分区重置偏移量,避免共享变量冲突 });
关键优化点
- 固定分区偏移量:预先计算每个分区在buffer中的位置,彻底消灭共享
counter带来的竞态问题。 - 分区级细粒度锁:每个分区对应独立的锁对象,只有操作当前分区时才会锁定,既保证了线程安全,又不会影响其他分区的并行效率。
- 安全资源管理:用
using语句自动管理文件流,避免资源泄漏。 - 原子操作封装:将单个分区的读buffer、写ADLS、清buffer操作放在同一个锁块里,确保这些步骤的原子性,不会被其他线程打断。
额外优化建议
如果你的buffer只是临时存储数据,其实可以让每个线程使用独立的临时缓冲区,完全避免共享buffer的同步问题,进一步提升并发效率(只要ADLS的ConcurrentAppend支持并行写入不同偏移量即可):
Parallel.For(0, partitionCount, i => { var filePath = @"C:\Users\t-chkum\Desktop\InputFiles\1MB\" + (i+1) + ".txt"; using (var stream = File.OpenRead(filePath)) { var tempBuffer = new byte[stream.Length]; stream.Read(tempBuffer, 0, (int)stream.Length); // 这里用预先计算好的每个文件的追加偏移量 var appendOffset = partitionOffsets[i]; client.ConcurrentAppend(fileName, true, tempBuffer, 0, (int)stream.Length, appendOffset); } });
这种方式不需要共享buffer和锁,每个线程独立处理自己的数据,只要偏移量计算正确就能安全并行写入。
内容的提问来源于stack exchange,提问作者kchahat20
相关产品推荐
相关产品推荐

