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

生产者/消费者代码未生成预期100万行数据,求排查问题

问题分析与修复方案

你的代码核心问题在于复用了同一个block数组对象,导致生产者在填充下一批数据时,直接覆盖了已经加入BlockingCollection的数组内容,消费者读取时拿到的是被修改后的数组,再加上你设置null作为终止符的逻辑,最终导致写入的行数远少于预期。

具体错误点拆解

  1. 数组引用复用
    你在生产者任务外部定义了String[] block = new string[arrSize];,每次填满数组后只是重置blockNum继续填充同一个数组,而buffer.Add(block)添加的是数组的引用,不是副本。当生产者继续修改这个数组时,消费者正在读取的还是同一个数组,数据被实时覆盖。

  2. 不必要的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 09:05:31