生产者/消费者代码未生成预期100万行数据,求排查问题
问题分析与修复方案
你的代码核心问题在于复用了同一个block数组对象,导致生产者在填充下一批数据时,直接覆盖了已经加入BlockingCollection的数组内容,消费者读取时拿到的是被修改后的数组,再加上你设置null作为终止符的逻辑,最终导致写入的行数远少于预期。
具体错误点拆解
数组引用复用
你在生产者任务外部定义了String[] block = new string[arrSize];,每次填满数组后只是重置blockNum继续填充同一个数组,而buffer.Add(block)添加的是数组的引用,不是副本。当生产者继续修改这个数组时,消费者正在读取的还是同一个数组,数据被实时覆盖。不必要的
null终止符逻辑
在赋值每个元素后,你设置了block[blockNum+1] = null,这会导致生产者在填充下一个元素前,提前把下一个位置设为null。如果消费者刚好在这个时候读取到这个null,就会直接终止当前block的遍历,跳过后面所有有效数据。
修复后的代码
我调整了生产者逻辑,每次填满一个block就创建新的数组,同时移除了多余的null设置,改用精确的元素数量处理最后一个不满的block:
class Program { static void Main(string[] args) { Random random = new Random(); String filename = @"d:\test_out"; long numlines = 1000000; var buffer = new BlockingCollection<string[]>(10); // 限制队列大小避免OOM int arrSize = 100; // 每个buffer块的字符串数量 Task producer = Task.Factory.StartNew(() => { long blockNum = 0; String[] block = new string[arrSize]; for (long i = 0; i < numlines; i++) { block[blockNum] = random.Next().ToString(); blockNum++; // 当块填满时,加入队列并创建新块 if (blockNum == arrSize) { buffer.Add(block); block = new string[arrSize]; blockNum = 0; } }; // 处理最后一个未填满的块 if (blockNum > 0) { // 创建仅包含有效元素的数组,避免null String[] finalBlock = new string[blockNum]; Array.Copy(block, finalBlock, blockNum); buffer.Add(finalBlock); } buffer.CompleteAdding(); }, TaskCreationOptions.LongRunning); Task consumer = Task.Factory.StartNew(() => { using (var outputFile = new StreamWriter(filename)) { foreach (string[] chunk in buffer.GetConsumingEnumerable()) { foreach (string value in chunk) { outputFile.WriteLine(value); } } } }, TaskCreationOptions.LongRunning); Task.WaitAll(producer, consumer); } }
修复效果说明
- 每次填满一个block就创建新数组,确保加入队列的数组不会被后续操作覆盖,消费者能读取到完整的原始数据。
- 最后一个block仅保留有效元素,消费者无需判断
null,直接遍历所有元素即可,保证总行数精确等于numlines设置的100万。
内容的提问来源于stack exchange,提问作者Sheyko Dmitriy
相关产品推荐
相关产品推荐

