将Polly与TPL数据流一起使用

Tod*_*ier 5 c# tpl-dataflow polly

数据处理管道和瞬态故障处理似乎是相辅相成的,所以我很想看看能否获得两个最佳的库,分别是TPL Dataflow和Polly,才能很好地协同工作。

首先,我想将错误处理策略应用于ActionBlock。理想情况下,我想将其封装在具有如下签名的块创建方法中:

ITargetBlock<T> CreatePollyBlock<T>(
    Action<T> act, ExecutionDataflowBlockOptions opts, Polly.Policy policy)
Run Code Online (Sandbox Code Playgroud)

从中简单地policy.Execute进行操作就足够容易了ActionBlock,但是我有以下两个要求:

  1. 在重试的情况下,我不想重试某个项目以使其优先于其他排队的项目。换句话说,当您失败时,您将转到行尾。
  2. 更重要的是,如果重试之前有一个等待期,我不想阻止新项目进入。如果ExecutionDataflowBlockOptions.MaxDegreeOfParallelism已设置,我不希望等待重试的项目根据该最大值进行“计数”。

为了满足这些要求,我想我需要一个“内部” ActionBlock与用户提供的ExecutionDataflowBlockOptions应用,并且一些“外部”块帖项内部块,并应用任何等待和重试逻辑(或任何策略规定)外内部块的上下文。这是我的第一次尝试:

// wrapper that provides a data item with mechanism to await completion
public class WorkItem<T>
{
    private readonly TaskCompletionSource<byte> _tcs = new TaskCompletionSource<byte>();

    public T Data { get; set; }
    public Task Completion => _tcs.Task;

    public void SetCompleted() => _tcs.SetResult(0);
    public void SetFailed(Exception ex) => _tcs.SetException(ex);
}

ITargetBlock<T> CreatePollyBlock<T>(Action<T> act, Policy policy, ExecutionDataflowBlockOptions opts) {
    // create a block that marks WorkItems completed, and allows
    // items to fault without faulting the entire block.
    var innerBlock = new ActionBlock<WorkItem<T>>(wi => {
        try {
            act(wi.Data);
            wi.SetCompleted();
        }
        catch (Exception ex) {
            wi.SetFailed(ex);
        }
    }, opts);

    return new ActionBlock<T>(async x => {
        await policy.ExecuteAsync(async () => {
            var workItem = new WorkItem<T> { Data = x };
            await innerBlock.SendAsync(workItem);
            await workItem.Completion;
        });
    });
}
Run Code Online (Sandbox Code Playgroud)

为了测试它,我创建了一个带有等待重试策略和虚拟方法的块,该方法在调用它的前3次(应用程序范围内)中引发异常。然后我给它一些数据:

"a", "b", "c", "d", "e", "f"
Run Code Online (Sandbox Code Playgroud)

我希望a,b和c失败并走到最后。但是我观察到它们按以下顺序击中了内部块的动作:

"a", "a", "a", "a", "b", "c", "d", "e", "f"
Run Code Online (Sandbox Code Playgroud)

本质上,我不能满足自己的要求,而且很容易理解为什么:外部块直到当前项目的所有重试都发生后才允许新项目进入。一个简单但看似骇人听闻的解决方案是MaxDegreeOfParallelism在外部块中增加很大的价值:

return new ActionBlock<T>(async x => {
    await policy.ExecuteAsync(async () => {
        var workItem = new WorkItem<T> { Data = x };
        await innerBlock.SendAsync(workItem);
        await workItem.Completion;
    });
}, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 100 });
Run Code Online (Sandbox Code Playgroud)

通过这种更改,我观察到实际上有新项目在重试之前进入,但是我也释放了一些混乱。进入内部块已成为随机的,尽管前三个项始终在末尾:

"a", "e", "b", "c", "d", "a", "e", "b"
Run Code Online (Sandbox Code Playgroud)

所以这好一点。但理想情况下,我希望保留订单:

"a", "b", "c", "d", "e", "a", "b", "c"
Run Code Online (Sandbox Code Playgroud)

这就是我遇到的问题,对此进行推理,我想知道在这些约束下是否有可能,特别是的内部CreatePollyBlock可以执行策略但不能定义策略。例如,如果这些内部组件可以提供重试lambda,我认为这将为我提供更多选择。但这是策略定义的一部分,按照这种设计,我无法做到这一点。

在此先感谢您的帮助。