You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.07 21:37:38