批量迁移中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

