【问题标题】:Alternate to Dataflow BroadcastBlock with guaranteed delivery替代 Dataflow BroadcastBlock 并保证交付
【发布时间】:2014-08-01 14:04:29
【问题描述】:

我需要某种对象,其行为类似于广播块,但有保证的交付。所以我使用了来自this question 的答案。但我并不太清楚这里的执行流程。我有一个控制台应用程序。这是我的代码:

static void Main(string[] args)
{
    ExecutionDataflowBlockOptions execopt = new ExecutionDataflowBlockOptions { BoundedCapacity = 5 };
    List<ActionBlock<int>> blocks = new List<ActionBlock<int>>();

    for (int i = 0; i <= 10; i++)
        blocks.Add(new ActionBlock<int>(num => 
        {
            int coef = i;
            Console.WriteLine(Thread.CurrentThread.ManagedThreadId + ". " + num * coef); 
        }, execopt));

    ActionBlock<int> broadcaster = new ActionBlock<int>(async num => 
    {
        foreach (ActionBlock<int> block in blocks) await block.SendAsync(num);
    }, execopt);

    broadcaster.Completion.ContinueWith(task =>
        {
            foreach (ActionBlock<int> block in blocks) block.Complete();
        });

    Task producer = Produce(broadcaster);
    List<Task> ToWait = new List<Task>();
    foreach (ActionBlock<int> block in blocks) ToWait.Add(block.Completion);
    ToWait.Add(producer);

    Task.WaitAll(ToWait.ToArray());

    Console.ReadLine();
}

static async Task Produce(ActionBlock<int> broadcaster)
{
    for (int i = 0; i <= 15; i++) await broadcaster.SendAsync(i);

    broadcaster.Complete();
}

每个数字都必须按顺序处理,所以我不能在广播块中使用 MaxDegreeOfParallelism。但是所有接收到该数字的动作块都可以并行运行。

那么问题来了:

在输出中我可以看到不同的线程 ID。我是否正确理解它的工作原理如下:

在广播公司中执行await block.SendAsync(num);。 如果当前块尚未准备好接受该数字,则执行退出广播器并在 Task.WaitAll 处挂起。 当 block 接受数字时,广播器中的 foreach 语句的其余部分在线程池中执行。 直到最后都一样。 foreach 的每次迭代都在线程池中执行。但实际上它是按顺序发生的。

我的理解是对还是错? 如何更改此代码以异步将号码发送到所有块?

为了确保如果其中一个块目前还没有准备好接收号码,我不会等待它,所有其他准备好的人都会收到号码。并且所有块都可以并行运行。并保证交货。

【问题讨论】:

    标签: c# multithreading task-parallel-library async-await tpl-dataflow


    【解决方案1】:

    假设您想通过 broadcaster 一次处理一个项目,同时使目标块能够同时接收该项目,您需要更改 broadcaster 以同时向所有块提供数量,然后异步等待所有人一起接受,然后再转到下一个号码:

    var broadcaster = new ActionBlock<int>(async num => 
    {
        var tasks = new List<Task>();
        foreach (var block in blocks)
        {
            tasks.Add(block.SendAsync(num));
        }
        await Task.WhenAll(tasks);
    }, execopt);
    

    现在,在这种情况下,如果您在等待后没有工作,您可以稍微优化,同时仍然返回等待的任务:

    ActionBlock<int> broadcaster = new ActionBlock<int>(
        num => Task.WhenAll(blocks.Select(block => block.SendAsync(num))), execopt);
    

    【讨论】:

    • 是否可以不等待所有解析器完成后再移动到下一个数字?我的意思是,只要有一些解析器的缓冲区可供接收,广播者就会发送给它。这样我就不会等待最慢的。我想我在问题中做了一些错误的解释。我需要每个解析器来处理所有数字,以便接收它们。但我不需要在每个数字之后等待所有解析器完成。
    • @ПавелБирюков 当您调用SendAsync 并且目标缓冲区中有空间时,返回的任务会立即完成。您只会在其中一个缓冲区已满时等待。您可以增加该缓冲区,但如果不确保项目移动到下一个块,我不会继续。
    • 但是当其中一个缓冲区已满时,我可以不等待,而是继续下一个数字以清空缓冲区吗?
    • 可以,但这可能会导致内存泄漏。您需要存储额外的数字以便以后能够将它们发送到完整块,这就像创建另一个您无法控制的缓冲区。如果内存不是问题,您可以简单地增加块的有限容量(甚至使其无限容量),但这可能会填满您的 RAM。
    • @MarcL。 TPL Dataflow 是一个内存库。如果您需要数据在崩溃中幸存下来,您当然不能依赖它,您应该使用持久排队机制。顺便说一句,由于 .NET 是托管的,只要您不使用不安全的代码,就不会发生缓冲区溢出(但是您可能会因堆栈溢出而崩溃)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-03-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-01-28
    • 2014-10-23
    • 2014-12-25
    相关资源
    最近更新 更多