5

我已经编写了一个示例测试来复制该问题。这不是我的实际代码,我试图写一个小复制。如果将边界容量增加到迭代次数,有效地使其没有边界,则不会死锁,如果将最大并行度设置为像 1 这样的小数,则不会死锁。

同样,我知道下面的代码不是很好,但我实际发现的代码要大得多且难以理解。基本上有一个与远程资源的连接的阻塞对象池,并且流中的几个块使用了该连接。

关于如何解决这个问题的任何想法?乍一看,这似乎是数据流的问题。当我停下来查看线程时,我看到许多线程在 Add 时被阻塞,0 个线程在 take 时被阻塞。addBlocks 出站队列中有几个项目尚未传播到 takeblock,因此它被卡住或死锁。

    var blockingCollection = new BlockingCollection<int>(10000);

    var takeBlock = new ActionBlock<int>((i) =>
    {
        int j = blockingCollection.Take();

    }, new ExecutionDataflowBlockOptions()
           {
              MaxDegreeOfParallelism = 20,
              SingleProducerConstrained = true
           });

    var addBlock = new TransformBlock<int, int>((i) => 
    {
        blockingCollection.Add(i);
        return i;

    }, new ExecutionDataflowBlockOptions()
           {
              MaxDegreeOfParallelism = 20
           });

    addBlock.LinkTo(takeBlock, new DataflowLinkOptions()
          {
             PropagateCompletion = true
          });

    for (int i = 0; i < 100000; i++)
    {
        addBlock.Post(i);
    }

    addBlock.Complete();
    await addBlock.Completion;
    await takeBlock.Completion;
4

1 回答 1

3

TPL Dataflow 不打算与阻塞很多的代码一起使用,我认为这个问题源于此。

我无法弄清楚到底发生了什么,但我认为解决方案是使用非阻塞集合。方便的是,Dataflow 以BufferBlock. 这样,您的代码将如下所示:

var bufferBlock = new BufferBlock<int>(
    new DataflowBlockOptions { BoundedCapacity = 10000 });

var takeBlock = new ActionBlock<int>(
    async i =>
    {
        int j = await bufferBlock.ReceiveAsync();
    }, new ExecutionDataflowBlockOptions
    {
        MaxDegreeOfParallelism = 20,
        SingleProducerConstrained = true
    });

var addBlock = new TransformBlock<int, int>(
    async i =>
    {
        await bufferBlock.SendAsync(i);
        return i;
    }, new ExecutionDataflowBlockOptions
    {
        MaxDegreeOfParallelism = 20
    });

尽管我发现您的代码的整个设计都很可疑。如果您想与块的正常结果一起发送一些附加数据,请将该块的输出类型更改为包含该附加数据的类型。

于 2014-05-28T17:45:07.413 回答