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

批量迁移中Blocking Collection出现重复条目问题求助

问题分析与修复方案

首先,你遇到的重复条目问题,核心原因是对象实例化的位置错误,导致多个线程复用同一个MigrationObject引用,最终属性被覆盖,看起来像是重复数据。我们一步步拆解问题并修复:

1. 最关键的重复原因:对象复用

在内部的Parallel.ForEach(Partitioner.Create(...))代码块中,你把MigrationObject migSource = new MigrationObject();放在了range循环的外面——这意味着所有的循环迭代(不管是哪个range、哪个线程)都在修改同一个对象的属性,然后把同一个引用添加到sourceCollection里。当多个线程同时修改这个对象时,后面的修改会覆盖前面的,最终你从集合里读取到的都是最后几次修改的值,自然就出现了“重复”条目。

修复方法:把MigrationObject的实例化放到内部的for循环里面,确保每一条记录都对应一个全新的对象实例:

Parallel.ForEach(Partitioner.Create(0, dt.Rows.Count), (range, state) => {
    for (int i = range.Item1; i < range.Item2; i++)
    {
        // 每一条记录都创建新的实例
        MigrationObject migSource = new MigrationObject();
        migSource.PAN = dt.Rows[i]["CUST_ID"].ToString();
        migSource.PinOffset = dt.Rows[i]["CODE"].ToString();
        migSource.PinBlockNew = dt.Rows[i]["BLOCK_NEW"].ToString();
        migSource.RelationshipNum = dt.Rows[i]["RELATIONSHIP_NUM"].ToString();
        Console.WriteLine(@"PAN " + migSource.PAN + " Rel " + migSource.RelationshipNum + " for ranges : " + range.Item1 + " TO " + range.Item2);
        sourceCollection.TryAdd(migSource);
    }
});

2. 分页逻辑的边界漏洞

当前计算setsCount的方式是Convert.ToInt32(cmd.ExecuteScalar()) / recordsInSet,这会导致当总记录数不是recordsInSet的整数倍时,最后一部分数据会被漏掉(比如总记录101条,会只处理前100条)。

修复方法:用向上取整的方式计算总批次,同时处理最后一批的边界:

int totalRecords = Convert.ToInt32(cmd.ExecuteScalar());
setsCount = (totalRecords + recordsInSet - 1) / recordsInSet; // 向上取整

// 循环生成批次时,修正最后一批的结束值
int rangeFrom = 1;
int rangeTo = recordsInSet;
for (int i = 0; i < setsCount; i++)
{
    if (rangeTo > totalRecords)
    {
        rangeTo = totalRecords;
    }
    chunks.Add(i, new Tuple<int, int>(rangeFrom, rangeTo));
    rangeFrom = rangeTo + 1;
    rangeTo = rangeTo + recordsInSet;
}

3. SQL语句的语法错误与安全问题

你当前的SQL语句直接把chunk.Value.Item1和chunk.Value.Item2写在字符串里,这会导致SQL引擎把它们当成列名,直接执行会报语法错误。同时,这种写法存在SQL注入风险,必须用参数化查询。

修复方法:修改SQL查询为参数化形式:

string command = @"SELECT * FROM ( 
    SELECT RELATIONSHIP_NUM, CUST_ID, CODE, BLOCK_NEW, 
    ROW_NUMBER() over ( order by RELATIONSHIP_NUM, CUST_ID) as RowNum 
    FROM MyTable) SUB 
WHERE SUB.RowNum BETWEEN @Start AND @End";

SqlDataAdapter da = new SqlDataAdapter(command, localConn);
da.SelectCommand.Parameters.AddWithValue("@Start", chunk.Value.Item1);
da.SelectCommand.Parameters.AddWithValue("@End", chunk.Value.Item2);

4. 资源释放的冗余操作

你在try/finally里手动关闭和释放localConn,但已经用了using (localConn),using块会自动处理连接的关闭和释放,所以finally里的代码是冗余的,可以直接删掉。

5. 并行度控制(针对百万级数据的优化)

处理150万条记录时,无限制的Parallel.ForEach会创建大量数据库连接,可能耗尽连接池导致性能下降甚至报错。建议限制并行度:

Parallel.ForEach(chunks, new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount }, (chunk) => {
    // 内部代码不变
});

修复后的完整代码片段

public async Task BulkMigrationAsync(string clearPVK, string EncZPK) {
    BlockingCollection<MigrationObject> sourceCollection = new BlockingCollection<MigrationObject>();
    int recordsInSet = 25;
    Dictionary<int, Tuple<int, int>> chunks = new Dictionary<int, Tuple<int, int>>();
    
    using (SqlConnection conn = new SqlConnection(_Database.ConnectionString))
    {
        conn.Open();
        SqlCommand cmd = new SqlCommand("SELECT COUNT(*) FROM MyTable", conn);
        int totalRecords = Convert.ToInt32(cmd.ExecuteScalar());
        int setsCount = (totalRecords + recordsInSet - 1) / recordsInSet; // 向上取整
        
        int rangeFrom = 1;
        int rangeTo = recordsInSet;
        for (int i = 0; i < setsCount; i++)
        {
            if (rangeTo > totalRecords)
            {
                rangeTo = totalRecords;
            }
            chunks.Add(i, new Tuple<int, int>(rangeFrom, rangeTo));
            rangeFrom = rangeTo + 1;
            rangeTo = rangeTo + recordsInSet;
        }
    }

    Task taskFetch = Task.Factory.StartNew(() => {
        // 限制并行度,避免连接池耗尽
        Parallel.ForEach(chunks, new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount }, (chunk) => {
            DataTable dt = new DataTable();
            string command = @"SELECT * FROM ( 
                SELECT RELATIONSHIP_NUM, CUST_ID, CODE, BLOCK_NEW, 
                ROW_NUMBER() over ( order by RELATIONSHIP_NUM, CUST_ID) as RowNum 
                FROM MyTable) SUB 
            WHERE SUB.RowNum BETWEEN @Start AND @End";
            
            using (SqlConnection localConn = new SqlConnection(_Database.ConnectionString))
            using (SqlDataAdapter da = new SqlDataAdapter(command, localConn))
            {
                da.SelectCommand.Parameters.AddWithValue("@Start", chunk.Value.Item1);
                da.SelectCommand.Parameters.AddWithValue("@End", chunk.Value.Item2);
                da.Fill(dt);
            }

            Parallel.ForEach(Partitioner.Create(0, dt.Rows.Count), (range, state) => {
                for (int i = range.Item1; i < range.Item2; i++)
                {
                    MigrationObject migSource = new MigrationObject();
                    migSource.PAN = dt.Rows[i]["CUST_ID"].ToString();
                    migSource.PinOffset = dt.Rows[i]["CODE"].ToString();
                    migSource.PinBlockNew = dt.Rows[i]["BLOCK_NEW"].ToString();
                    migSource.RelationshipNum = dt.Rows[i]["RELATIONSHIP_NUM"].ToString();
                    Console.WriteLine(@"PAN " + migSource.PAN + " Rel " + migSource.RelationshipNum + " for ranges : " + range.Item1 + " TO " + range.Item2);
                    sourceCollection.TryAdd(migSource);
                }
            });
        });
    });

    await taskFetch;
    sourceCollection.CompleteAdding();

    while (!sourceCollection.IsCompleted)
    {
        if (sourceCollection.TryTake(out MigrationObject mig))
        {
            // 去掉不必要的Task.Delay,除非有特殊的速率控制需求
            // await Task.Delay(50);
            Console.WriteLine(" Rel " + mig.RelationshipNum + " PAN " + mig.PAN);
        }
    }
}

最后补充一点:如果你的业务场景是百万级数据迁移,BlockingCollection虽然可用,但也可以考虑更高效的管道模式(比如TPL Dataflow),不过先解决当前的重复问题是首要的。

内容的提问来源于stack exchange,提问作者Sadiq

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 03:32:57