【问题标题】:Handle exceptions with TPL Dataflow blocks使用 TPL 数据流块处理异常
【发布时间】:2019-07-09 09:27:12
【问题描述】:

我有一个简单的 tpl 数据流,它基本上可以完成一些任务。 我注意到当任何数据块中出现异常时,它并没有被初始父块调用者捕获。 我添加了一些手动代码来检查异常,但似乎不是正确的方法。

if (readBlock.Completion.Exception != null
    || saveBlockJoinedProcess.Completion.Exception != null
    || processBlock1.Completion.Exception != null
    || processBlock2.Completion.Exception != null)
{
    throw readBlock.Completion.Exception;
}

我在网上查看了建议的方法,但没有看到任何明显的东西。 所以我在下面创建了一些示例代码,并希望获得一些关于更好解决方案的指导:

using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using System.Threading.Tasks.Dataflow;

namespace TPLDataflow
{
    class Program
    {
        static void Main(string[] args)
        {
            try
            {
                //ProcessB();
                ProcessA();
            }
            catch (Exception e)
            {
                Console.WriteLine("Exception in Process!");
                throw new Exception($"exception:{e}");
            }
            Console.WriteLine("Processing complete!");
            Console.ReadLine();
        }

        private static void ProcessB()
        {
            Task.WhenAll(Task.Run(() => DoSomething(1, "ProcessB"))).Wait();
        }

        private static void ProcessA()
        {
            var random = new Random();
            var readBlock = new TransformBlock<int, int>(x =>
            {
                try { return DoSomething(x, "readBlock"); }
                catch (Exception e) { throw e; }
            }); //1

            var braodcastBlock = new BroadcastBlock<int>(i => i); // ⬅ Here

            var processBlock1 = new TransformBlock<int, int>(x =>
                DoSomethingAsync(5, "processBlock1")); //2
            var processBlock2 = new TransformBlock<int, int>(x =>
                DoSomethingAsync(2, "processBlock2")); //3

            //var saveBlock =
            //    new ActionBlock<int>(
            //    x => Save(x)); //4

            var saveBlockJoinedProcess =
                new ActionBlock<Tuple<int, int>>(
                x => SaveJoined(x.Item1, x.Item2)); //4

            var saveBlockJoin = new JoinBlock<int, int>();

            readBlock.LinkTo(braodcastBlock, new DataflowLinkOptions
                { PropagateCompletion = true });

            braodcastBlock.LinkTo(processBlock1,
                new DataflowLinkOptions { PropagateCompletion = true }); //5

            braodcastBlock.LinkTo(processBlock2,
                new DataflowLinkOptions { PropagateCompletion = true }); //6


            processBlock1.LinkTo(
                saveBlockJoin.Target1); //7

            processBlock2.LinkTo(
                saveBlockJoin.Target2); //8

            saveBlockJoin.LinkTo(saveBlockJoinedProcess,
                new DataflowLinkOptions { PropagateCompletion = true });

            readBlock.Post(1); //10
                               //readBlock.Post(2); //10

            Task.WhenAll(processBlock1.Completion,processBlock2.Completion)
                .ContinueWith(_ => saveBlockJoin.Complete());

            readBlock.Complete(); //12
            saveBlockJoinedProcess.Completion.Wait(); //13
            if (readBlock.Completion.Exception != null
                || saveBlockJoinedProcess.Completion.Exception != null
                || processBlock1.Completion.Exception != null
                || processBlock2.Completion.Exception != null)
            {
                throw readBlock.Completion.Exception;
            }
        }
        private static int DoSomething(int i, string method)
        {
            Console.WriteLine($"Do Something, callng method : { method}");
            throw new Exception("Fake Exception!");
            return i;
        }
        private static async Task<int> DoSomethingAsync(int i, string method)
        {
            Console.WriteLine($"Do SomethingAsync");
            throw new Exception("Fake Exception!");
            await Task.Delay(new TimeSpan(0, 0, i));
            Console.WriteLine($"Do Something : {i}, callng method : { method}");
            return i;
        }
        private static void Save(int x)
        {

            Console.WriteLine("Save!");
        }
        private static void SaveJoined(int x, int y)
        {
            Thread.Sleep(new TimeSpan(0, 0, 10));
            Console.WriteLine("Save Joined!");
        }
    }
}

【问题讨论】:

  • 如果您愿意接受替代方案,我建议您使用 Observables 采用 Rx 方式 - 它们更容易让您绕开脑袋
  • @KrzysztofSkowronek 我也喜欢 Rx,但使用 TPL DF 作为尝试它来完成这个特定任务的一种方式 :)

标签: c# .net task-parallel-library tpl-dataflow


【解决方案1】:

我在网上查看了建议的方法,但没有看到任何明显的内容。

如果你有管道(或多或少),那么常用的方法是使用PropagateCompletion 关闭管道。如果您有更复杂的拓扑,则需要手动完成块。

在您的情况下,您在此处尝试传播:

Task.WhenAll(
    processBlock1.Completion,
    processBlock2.Completion)
    .ContinueWith(_ => saveBlockJoin.Complete());

但此代码不会传播异常。当processBlock1.CompletionprocessBlock2.Completion 都完成时,saveBlockJoin 完成成功

更好的解决方案是使用await 而不是ContinueWith

async Task PropagateToSaveBlockJoin()
{
    try
    {
        await Task.WhenAll(processBlock1.Completion, processBlock2.Completion);
        saveBlockJoin.Complete();
    }
    catch (Exception ex)
    {
        ((IDataflowBlock)saveBlockJoin).Fault(ex);
    }
}
_ = PropagateToSaveBlockJoin();

使用await 鼓励您处理异常,您可以通过将它们传递给Fault 来传播异常。

【讨论】:

    【解决方案2】:

    开箱即用的 TPL 数据流不支持在管道中向后传播错误,这在块具有有限容量时尤其令人讨厌。在这种情况下,下游块中的错误可能会导致其前面的块无限期地阻塞。我知道的唯一解决方案是使用取消功能,并在有人失败时取消所有块。这是如何做到的。首先创建一个CancellationTokenSource

    var cts = new CancellationTokenSource();
    

    然后一个一个地创建块,在所有的选项中嵌入相同的CancellationToken

    var options = new ExecutionDataflowBlockOptions()
        { BoundedCapacity = 10, CancellationToken = cts.Token };
    
    var block1 = new TransformBlock<double, double>(Math.Sqrt, options);
    var block2 = new ActionBlock<double>(Console.WriteLine, options);
    

    然后将块链接在一起,包括PropagateCompletion 设置:

    block1.LinkTo(block2, new DataflowLinkOptions { PropagateCompletion = true });
    

    最后使用扩展方法触发CancellationTokenSource在异常情况下的取消:

    block1.OnFaultedCancel(cts);
    block2.OnFaultedCancel(cts);
    

    OnFaultedCancel扩展方法如下图:

    public static class DataflowExtensions
    {
        public static void OnFaultedCancel(this IDataflowBlock dataflowBlock,
            CancellationTokenSource cts)
        {
            dataflowBlock.Completion.ContinueWith(_ => cts.Cancel(), default,
                TaskContinuationOptions.OnlyOnFaulted |
                TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Default);
        }
    }
    

    【讨论】:

    • 这样,感觉很自然。
    【解决方案3】:

    乍一看,如果只有一些小点(不看你的架构)。在我看来,您混合了一些新的和一些旧的结构。还有一些代码部分是不必要的。

    例如:

    private static void ProcessB()
    {
        Task.WhenAll(Task.Run(() => DoSomething(1, "ProcessB"))).Wait();
    }
    

    使用 Wait() 方法,如果发生任何异常,它们将被包装在 System.AggregateException 中。在我看来,这样更好:

    private static async Task ProcessBAsync()
    {
        await Task.Run(() => DoSomething(1, "ProcessB"));
    }
    

    使用 async-await,如果发生异常,await 语句会重新抛出包装在 System.AggregateException 中的第一个异常。这使您可以尝试捕获具体的异常类型并仅处理您真正可以处理的情况。

    另一件事是这部分代码:

    private static void ProcessA()
            {
                var random = new Random();
                var readBlock = new TransformBlock<int, int>(
                        x => 
                        { 
                        try { return DoSomething(x, "readBlock"); } 
                        catch (Exception e) 
                        { 
                        throw e; 
                        } 
                        },
                        new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 1 }); //1
    

    为什么捕获异常只是为了重新抛出它?在这种情况下,try-catch 是多余的。

    这里是:

    private static void SaveJoined(int x, int y)
    {
        Thread.Sleep(new TimeSpan(0, 0, 10));
        Console.WriteLine("Save Joined!");
    }
    

    最好使用 await Task.Delay(....)。使用Task.Delay(...),您的应用程序不会冻结。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-02-16
      • 1970-01-01
      • 1970-01-01
      • 2012-04-10
      • 2015-12-22
      • 2012-04-14
      • 1970-01-01
      • 2012-02-01
      相关资源
      最近更新 更多