【问题标题】:TPL Dataflow, BroadcastBlock to BatchBlocksTPL 数据流,BroadcastBlock 到 BatchBlocks
【发布时间】:2015-04-25 05:04:57
【问题描述】:

我在将BroadcastBlock(s) 连接到BatchBlocks 时遇到问题。场景是来源是BroadcastBlocks,收件人是BatchBlocks

在下面的简化代码中,只有一个补充动作块执行。我什至将每个 BatchBlock 的 batchSize 设置为 1 来说明问题。

将 Greedy 设置为“true”将使 2 ActionBlocks 执行,但这不是我想要的,因为它会导致 BatchBlock 继续进行,即使它尚未完成。有什么想法吗?

class Program
{
    static void Main(string[] args)
    {
        // My possible sources are BroadcastBlocks. Could be more
        var source1 = new BroadcastBlock<int>(z => z);

        // batch 1
        // can be many potential sources, one for now
        // I want all sources to arrive first before proceeding
        var batch1 = new BatchBlock<int>(1, new GroupingDataflowBlockOptions() { Greedy = false }); 
        var batch1Action = new ActionBlock<int[]>(arr =>
        {
            // this does not run sometimes
            Console.WriteLine("Received from batch 1 block!");
            foreach (var item in arr)
            {
                Console.WriteLine("Received {0}", item);
            }
        });

        batch1.LinkTo(batch1Action, new DataflowLinkOptions() { PropagateCompletion = true });

        // batch 2
        // can be many potential sources, one for now
        // I want all sources to arrive first before proceeding
        var batch2 = new BatchBlock<int>(1, new GroupingDataflowBlockOptions() { Greedy = false  });
        var batch2Action = new ActionBlock<int[]>(arr =>
        {
            // this does not run sometimes
            Console.WriteLine("Received from batch 2 block!");
            foreach (var item in arr)
            {
                Console.WriteLine("Received {0}", item);
            }
        });
        batch2.LinkTo(batch2Action, new DataflowLinkOptions() { PropagateCompletion = true });

        // connect source(s)
        source1.LinkTo(batch1, new DataflowLinkOptions() { PropagateCompletion = true });
        source1.LinkTo(batch2, new DataflowLinkOptions() { PropagateCompletion = true });

        // fire
        source1.SendAsync(3);

        Task.WaitAll(new Task[] { batch1Action.Completion, batch2Action.Completion }); ;

        Console.ReadLine();
    }
}

【问题讨论】:

  • 我认为将Greedy 设置为true 正确的解决方案。如果您担心的是,它不会导致创建较小的批次。

标签: c# concurrency task-parallel-library tpl-dataflow


【解决方案1】:

支持非贪心功能的 TPL Dataflow 库的内部机制似乎存在缺陷。发生的情况是BatchBlock 配置为非贪婪将Postpone 链接块提供的所有消息,而不是接受它们。它与已推迟的消息保持一个内部队列,当它们的数量达到其BatchSize 配置时,它会尝试使用推迟的消息,如果成功,它将按预期将它们传播到下游。问题在于,像BroadcastBlockBufferBlock 这样的源块将停止向已推迟先前提供的消息的块提供更多消息,直到它消耗了这条消息。这两种行为的结合导致了死锁。无法向前推进,因为BatchBlock 在消费推迟的消息之前等待提供更多消息,而BroadcastBlock 在提供更多消息之前等待推迟的消息被消费...

这种情况只发生在BatchSize 大于一的情况下(这是此块的典型配置)。

Here 是这个问题的一个演示。作为源,它使用更常见的BufferBlock 而不是BroadcastBlock。 10 条消息发布到三块管道,预期行为是消息通过管道流到最后一个块。实际上什么也没发生,所有消息都停留在第一个块中。

using System;
using System.Threading;
using System.Threading.Tasks.Dataflow;

public static class Program
{
    static void Main(string[] args)
    {
        var bufferBlock = new BufferBlock<int>();

        var batchBlock = new BatchBlock<int>(batchSize: 2,
            new GroupingDataflowBlockOptions() { Greedy = false });

        var actionBlock = new ActionBlock<int[]>(batch =>
            Console.WriteLine($"Received: {String.Join(", ", batch)}"));

        bufferBlock.LinkTo(batchBlock,
            new DataflowLinkOptions() { PropagateCompletion = true });

        batchBlock.LinkTo(actionBlock,
            new DataflowLinkOptions() { PropagateCompletion = true });

        for (int i = 1; i <= 10; i++)
        {
            var accepted = bufferBlock.Post(i);
            Console.WriteLine(
                $"bufferBlock.Post({i}) {(accepted ? "accepted" : "rejected")}");
            Thread.Sleep(100);
        }

        bufferBlock.Complete();
        actionBlock.Completion.Wait(millisecondsTimeout: 1000);
        Console.WriteLine();
        Console.WriteLine($"bufferBlock.Completion: {bufferBlock.Completion.Status}");
        Console.WriteLine($"batchBlock.Completion:  {batchBlock.Completion.Status}");
        Console.WriteLine($"actionBlock.Completion: {actionBlock.Completion.Status}");
        Console.WriteLine($"bufferBlock.Count: {bufferBlock.Count}");
    }
}

输出:

bufferBlock.Post(1) accepted
bufferBlock.Post(2) accepted
bufferBlock.Post(3) accepted
bufferBlock.Post(4) accepted
bufferBlock.Post(5) accepted
bufferBlock.Post(6) accepted
bufferBlock.Post(7) accepted
bufferBlock.Post(8) accepted
bufferBlock.Post(9) accepted
bufferBlock.Post(10) accepted

bufferBlock.Completion: WaitingForActivation
batchBlock.Completion:  WaitingForActivation
actionBlock.Completion: WaitingForActivation
bufferBlock.Count: 10

我的猜测是,内部的offer-consume-reserve-release 机制已经过调整,以最大限度地支持BoundedCapacity 功能,这对许多应用程序来说至关重要,而且很少使用Greedy = false 功能未经彻底测试。

好消息是,在您的情况下,您实际上并不需要将Greedy 设置为false。默认贪婪模式下的BatchBlock 不会传播比配置的BatchSize 更少的消息,除非它已被标记为已完成并传播任何剩余的消息,或者您在任意时刻手动调用其TriggerBatch 方法。非贪婪配置的预期用途是防止 complex graph scenarios 中的资源匮乏,块之间存在多个依赖关系。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多