【问题标题】:Microsoft TPL Dataflow - processing correlating requests synchronouslyMicrosoft TPL 数据流 - 同步处理相关请求
【发布时间】:2014-07-07 19:30:24
【问题描述】:

我提前为标题道歉,但这是我能想到的最好的描述动作。

要求是处理消息总线的请求。 进来的请求可能与关联或分组这些请求的 id 有关。 我想要的行为是让请求流同步处理相关的 id。 但是可以异步处理不同的 id。

我正在使用并发字典来跟踪正在处理的请求和链接中的谓词。

这是假设提供相关请求的同步处理。

但是我得到的行为是第一个请求得到处理,第二个请求被丢弃。

我已附加来自控制台应用程序的示例代码来模拟问题。

我们将不胜感激任何方向或反馈。

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

namespace ConsoleApplication2
{
    class Program
    {
        static void Main(string[] args)
        {
            var requestTracker = new ConcurrentDictionary<string, string>();

            var bufferBlock = new BufferBlock<Request>();

            var actionBlock = new ActionBlock<Request>(x => 
            {
                Console.WriteLine("processing item {0}",x.Name);
                Thread.Sleep(5000);
                string itemOut = null;
                requestTracker.TryRemove(x.Id, out itemOut);
            });

            bufferBlock.LinkTo(actionBlock, x => requestTracker.TryAdd(x.Id,x.Name));


            var publisher = Task.Run(() =>
            {
                var request = new Request("item_1", "first item");
                bufferBlock.SendAsync(request);

                var request_1 = new Request("item_1", "second item");
                bufferBlock.SendAsync(request_1);

            });

            publisher.Wait();
            Console.ReadLine();
        }
    }

    public class Request
    {
        public Request(string id, string name)
        {
            this.Id = id;
            this.Name = name;
        }
        public string Id { get; set; }
        public string Name { get; set; }
    }
}

【问题讨论】:

  • 您应该让异常通过数据流的管道传播,这样您就可以看到出了什么问题。在MSDN Walkthrough 的末尾查看 MS 的完整示例。然后,您可以处理 AggregateException 以找出问题所在。
  • 你的意思是你想让一个具有相同id的组一个接一个地处理,而组可以同时处理?如果是这样,你的答案是:stackoverflow.com/q/21010024/885318
  • @I3arnon - 您的解决方案似乎正是我正在寻找的。我在这里有点懒,但也许你有更多的细节。我猜你得到的密钥是动态的,即基本上消息的爆发将具有相同的密钥并且它一直在变化。当一个动作块忙于处理一条消息时,它会更新一个字典,说我正忙于这个请求,并且任何与该键匹配的后续请求都被委派给该动作块?我说的对吗?
  • @I3arnon 我认为这在这里不必要地复杂。
  • @svick 你有什么建议?

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


【解决方案1】:
  1. 你说你想并行处理一些请求(至少我假设这就是你所说的“异步”),但ActionBlock 默认情况下不是并行的。要更改它,请设置MaxDegreeOfParallelism。

  2. 您正在尝试使用 TryAdd() 作为过滤器,但这不起作用有两个原因:

    1. 过滤器只被调用一次,它不会自动重试或类似的东西。这意味着如果一个项目没有通过,它就永远不会通过,即使在阻止它的项目完成之后也是如此。
    2. 如果一个项目卡在一个块的输出队列中,则没有其他项目会离开该块。即使您以某种方式解决了上一个问题,这也可能会显着降低并行度。
  3. 我认为这里最简单的解决方案是为每个组设置一个块,这样,每个组的项目将按顺序处理,但不同组的项目将并行处理。在代码中,它可能看起来像:

    var processingBlocks = new Dictionary<string, ActionBlock<Request>>();
    
    var splitterBlock = new ActionBlock<Request>(request =>
    {
        ActionBlock<Request> processingBlock;
    
        if (!processingBlocks.TryGetValue(request.Id, out processingBlock))
        {
            processingBlock = processingBlocks[request.Id] =
                new ActionBlock<Request>(r => /* process the request here */);
        }
    
        processingBlock.Post(request);
    });
    

    这种方法的问题是组的处理块永远不会消失。如果您负担不起(这是内存泄漏),因为您将拥有大量组,那么hashing approach suggested by I3arnon 是您的最佳选择。

【讨论】:

  • 感谢 svick 的帮助 - 实际上我最终使用了 I3arnon 方法。对于我尝试做的事情来说,它更健壮,更不容易出错。
【解决方案2】:

我相信这是因为您的LinkTo() 设置不正确。通过拥有LinkTo() 并将函数作为参数传递,您正在添加条件。所以这一行:

bufferBlock.LinkTo(actionBlock, x => requestTracker.TryAdd(x.Id, x.Name));

本质上是说,如果您能够添加到并发字典中,则将数据从 bufferBlock 传递到 actionBlock,这不一定有意义(至少在您的示例代码中)

相反,您应该将您的 bufferBlock 链接到没有 lambda 的操作块,因为在这种情况下您不需要条件链接(至少根据您的示例代码我不这么认为)。

另外,看看这个 SO 问题,看看您是否应该使用 SendAsync() 或 Post(),因为 Post() 可以更容易处理,只需将数据添加到管道中:TPL Dataflow, whats the functional difference between Post() and SendAsync()?。 SendAsync 将返回一个任务,而 Post 将根据成功进入管道返回 true/false。

因此,要从根本上找出问题所在,您需要处理块的延续。在 MSDN 的 TPL 数据流介绍中有一个很好的教程:Create a DataFlow Pipeline 它基本上看起来像这样:

//link to section
bufferBlock.LinkTo(actionBlock);
//continuations
bufferBlock.Completion.ContinueWith(t =>
{
     if(t.IsFaulted)  ((IDataFlowBlock).actionBlock).Fault(t.Exception); //send the exception down the pipeline
     else actionBlock.Complete(); //tell the next block that we're done with the bufferblock
 });

然后您可以在等待管道时捕获异常 (AggregateException)。您是否真的需要在实际代码中使用并发字典进行跟踪,因为当它无法添加时,这可能会导致问题,因为当链接到谓词返回 false 时,它​​不会将数据传递到管道的下一个块.

【讨论】:

  • 感谢 JNYRanger,我会尝试一下。 TryAdd 在 TryGet 上的目的是为了线程安全。另一个是并发字典针对读取进行了优化。在我的情况下,2个相关的 ids 通过检查并且都试图添加。因此我选择了 TryAdd 在这种情况下,一个会赢得比赛条件,另一个会被搁置。
  • @rizan 是有道理的,但是为了让一切都通过,您需要第二个LinkTo() 来处理添加到 ConcurrentDictionary 的故障,否则您的管道将中断,这就是你正在经历。
  • 您的方法和您提供的样本在可以等待完成的环境中表现出色。我所做的是采用 I3arnon 的方法,并稍作调整。我首先使用 TDD 编写测试,所以我涵盖了所有基础。非常感谢您的帮助。
  • @rizan 很高兴我能提供帮助,至少可以为您指明正确的方向。我建议您回答自己的问题,向其他人展示您是如何解决问题的。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-06-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-07-28
  • 2018-10-11
  • 2021-10-21
相关资源
最近更新 更多