【问题标题】:Each element of InputQueue in TransformBlock<TInput, TOutput> is overwritten whenever a record is read每当读取记录时,TransformBlock<TInput, TOutput> 中 InputQueue 的每个元素都会被覆盖
【发布时间】:2020-11-11 17:05:36
【问题描述】:

目的

我正在尝试创建一个管道,在该管道中,我一次从文件中读取一个字节的记录,并将其发送到“缓冲区块”,后者将项目附加到缓冲区块中。这通过简单的 LinkTo () 方法链接到 TransformBlock ,后者将每个字节记录转换为 MyObject 对象。 以下是执行所有这些操作的整个方法:

    BufferBlock<byte[]> buffer = new BufferBlock<byte[]>();
    TransformBlock<byte[], MyObject> transform = new TransformBlock<byte[], MyObject>(bytes =>
    {
        return FromBytesTOMyObject(bytes);
    });

    private void ReadFileAndAppend()
    {
        buffer.LinkTo(transform, new DataflowLinkOptions() { PropagateCompletion = true });

        BinaryReader br = new BinaryReader(new FileStream("C:\\Users\\MyUser\\myFile.raw",FileMode.Open, FileAccess.Read, FileShare.Read));                                  
        int count;
        byte[] record = new byte[4000];

        // Post more messages to the block.
        Task.Run(async () =>
        {
            while ((count = br.Read(record, 0, record.Length)) != 0)
                await buffer.SendAsync(record);
            buffer.Complete();
        });
        transform.Completion.Wait();
        Console.WriteLine("");

在TransformBlock内部调用的方法下方:

static public MyObject FromBytesToMyObject(byte[] record)
    {
        MyObject object = new MyObject();
        object.counter = BitConverter.ToInt32(record, 0);
        object.nPoints = BitConverter.ToInt32(record, 4);

        for (int i = 0; i < object.nPoints; i++)
        {
            int index = i * 4;
            object.A[i] = BitConverter.ToSingle(record, index + 8);
        }
        return object;
    }

从FromBytesToMyObject()方法可以看出,每条读取的记录里面都有一个计数器。 所以每条记录永远不会有一个等于另一条记录的计数器(我还通过像 HxD 这样的字节读取器进行了检查)

问题

有了这个设置,我认为文件的读取和解释很顺利。但是在读取大约 50 条或更多记录后进入调试并在“while”中插入断点,我注意到在 TransformBlock 的 OutputQueue 中 具有相同计数器的记录组排队,因此记录相同。

例子:

仅考虑计数器的确切队列: 1,2,3,4,5,6,7,8,9,10。

我在 OutputQueue 中实际看到的队列:1,2,2,2,3,3,3,3,4,4,5,5 ....

你能解释一下我哪里错了吗?

【问题讨论】:

  • @EugeneSh。 IMO 它是 C# 。我已经编辑了标签
  • while (br.BaseStream.Position != br.BaseStream.Length) 行看起来很粗略。流读取循环的完成通常由返回零的reader.Read 方法确定。使用外部提供的byte[] 缓冲区看起来更加粗略。天知道这个缓冲区还能在哪里同时使用......
  • 没有minimal reproducible example 就不可能提供一个好的答案。也就是说,从描述来看,几乎可以肯定您对每个块都使用相同的引用类型对象实例,因此当然当您修改该对象时,对同一对象的所有其他引用都会显示相同的修改。
  • 附带说明,BufferBlock 可能是多余的。您可以直接输入TransformBlock,因为它有自己的内部输入队列。

标签: c# tpl-dataflow


【解决方案1】:

您一遍又一遍地重复使用相同的byte[] record,没有任何线程安全考虑。难怪事情进展得很快。如果要保证整个操作的正确性,每次都必须使用不同的缓冲区:

while (true)
{
    var record = new byte[4000];
    var count = binaryReader.Read(record, 0, record.Length);
    if (count == 0) break;
    Array.Resize(ref record, count);
    await buffer.SendAsync(record);
}
buffer.Complete();

如果您还关心性能并且不想让垃圾收集器负担过重,您应该查看ArrayPool 类。但要小心,因为这门课提供了新的方法来让自己一枪毙命。

【讨论】:

  • 如果每次都重新分配变量,Array.Resize(ref record, count); 是什么意思?
  • @NickMan 最后一个Read 很可能读取不到4000 个字节,因此在将缓冲区传递给TransformBlock 时,您需要某种方式来传达读取的确切字节数。您必须调整缓冲区的大小,或者将计数与缓冲区一起传递给块。
猜你喜欢
  • 2021-04-20
  • 1970-01-01
  • 2015-07-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2010-10-09
  • 2023-02-23
  • 2011-02-22
相关资源
最近更新 更多