如何用异常处理创建永无止境的DataFlow Mesh?

Ami*_*mit 4 c# concurrency task-parallel-library tpl-dataflow

我正在创建一个使用TPL DataFlow的任务处理器。我将遵循生产者消费者模型,在生产者模型中,生产者偶尔会生产一些要处理的商品,而消费者则一直在等待新商品的到来。这是我的代码:

async Task Main()
{
    var runner = new Runner();
    CancellationTokenSource cts = new CancellationTokenSource();
    Task runnerTask = runner.ExecuteAsync(cts.Token);

    await Task.WhenAll(runnerTask);
}

public class Runner
{
    public async Task ExecuteAsync(CancellationToken cancellationToken) {
        var random = new Random();

        ActionMeshProcessor processor = new ActionMeshProcessor();
        await processor.Init(cancellationToken);

        while (!cancellationToken.IsCancellationRequested)
        {
            await Task.Delay(TimeSpan.FromSeconds(1)); // wait before enqueuing more

            int[] items = GetItems(random.Next(3, 7));

            await processor.ProcessBlockAsync(items);
        }
    }

    private int[] GetItems(int count)
    {
        Random randNum = new Random();

        int[] arr = new int[count];
        for (int i = 0; i < count; i++)
        {
            arr[i] = randNum.Next(10, 20);
        }

        return arr;
    }
}

public class ActionMeshProcessor
{
    private TransformBlock<int, int> Transformer { get; set; }
    private ActionBlock<int> CompletionAnnouncer { get; set; }

    public async Task Init(CancellationToken cancellationToken)
    {
        var options = new ExecutionDataflowBlockOptions
        {
            CancellationToken = cancellationToken,
            MaxDegreeOfParallelism = 5,
            BoundedCapacity = 5
        };


        this.Transformer = new TransformBlock<int, int>(async input => {

            await Task.Delay(TimeSpan.FromSeconds(1)); //donig something complex here!

            if (input > 15)
            {
                throw new Exception($"I can't handle this number: {input}");
            }

            return input + 1;
        }, options);

        this.CompletionAnnouncer = new ActionBlock<int>(async input =>
        {
            Console.WriteLine($"Completed: {input}");

            await Task.FromResult(0);
        }, options);

        this.Transformer.LinkTo(this.CompletionAnnouncer);

        await Task.FromResult(0); // what do I await here?
    }

    public async Task ProcessBlockAsync(int[] arr)
    {
        foreach (var item in arr)
        {
            await this.Transformer.SendAsync(item); // await if there are no free slots
        }       
    }
}
Run Code Online (Sandbox Code Playgroud)

我在上面添加了条件检查,以抛出异常来模仿特殊情况。

这是我的问题:

  • 我可以在不降低整个网格的情况下处理上述网格中的异常的最佳方法是什么?

  • 有没有更好的方法来初始化/启动/继续永无止境的DataFlow网格?

  • 我在哪里等待完成?

我已经看过这个类似的问题

JSt*_*ard 5

例外情况

您的异步中没有任何东西,init它可以是标准的同步构造函数。您可以处理网格中的异常,而无需通过在提供给块的lamda中进行简单的try catch来降低网格的大小。然后,可以通过从网格中过滤结果或在以下块中忽略结果来处理这种情况。以下是过滤的示例。对于an的简单情况,int您可以使用int?和过滤掉以前的任何值,null或者,当然,您可以根据需要设置任何类型的魔术指示符值。如果您实际传递的是引用类型,则可以压入null或将数据项标记为脏数据,以使链接上的谓词可以检查该数据项。

public class ActionMeshProcessor {
    private TransformBlock<int, int?> Transformer { get; set; }
    private ActionBlock<int?> CompletionAnnouncer { get; set; }

    public ActionMeshProcessor(CancellationToken cancellationToken) {
        var options = new ExecutionDataflowBlockOptions {
            CancellationToken = cancellationToken,
            MaxDegreeOfParallelism = 5,
            BoundedCapacity = 5
        };


        this.Transformer = new TransformBlock<int, int?>(async input => {
            try {
                await Task.Delay(TimeSpan.FromSeconds(1)); //donig something complex here!

                if (input > 15) {
                    throw new Exception($"I can't handle this number: {input}");
                }

                return input + 1;
            } catch (Exception ex) {
                return null;
            }

        }, options);

        this.CompletionAnnouncer = new ActionBlock<int?>(async input =>
        {
            if (input == null) throw new ArgumentNullException("input");

            Console.WriteLine($"Completed: {input}");

            await Task.FromResult(0);
        }, options);

        //Filtering
        this.Transformer.LinkTo(this.CompletionAnnouncer, x => x != null);
        this.Transformer.LinkTo(DataflowBlock.NullTarget<int?>());
    }

    public async Task ProcessBlockAsync(int[] arr) {
        foreach (var item in arr) {
            await this.Transformer.SendAsync(item); // await if there are no free slots
        }
    }
}
Run Code Online (Sandbox Code Playgroud)

完成时间

您可以在关闭应用程序时公开Complete()和Completion从处理器中使用它们,并使用它们来await完成操作,假设那是您唯一一次关闭网格的情况。另外,请确保通过链接正确传播完成。

    //Filtering
    this.Transformer.LinkTo(this.CompletionAnnouncer, new DataflowLinkOptions() { PropagateCompletion = true }, x => x != null);
    this.Transformer.LinkTo(DataflowBlock.NullTarget<int?>());
}        

public void Complete() {
    Transformer.Complete();
}

public Task Completion {
    get { return CompletionAnnouncer.Completion; }
}
Run Code Online (Sandbox Code Playgroud)

然后,根据您的样本,最有可能完成工作的地方不在驱动您处理的循环之外:

public async Task ExecuteAsync(CancellationToken cancellationToken) {
    var random = new Random();

    ActionMeshProcessor processor = new ActionMeshProcessor();
    await processor.Init(cancellationToken);

    while (!cancellationToken.IsCancellationRequested) {
        await Task.Delay(TimeSpan.FromSeconds(1)); // wait before enqueuing more

        int[] items = GetItems(random.Next(3, 7));

        await processor.ProcessBlockAsync(items);
    }
    //asuming you don't intend to throw from cancellation
    processor.Complete();
    await processor.Completion();

}
Run Code Online (Sandbox Code Playgroud)