【问题标题】:TPL DataFlow with Lazy Source / stream of data具有惰性源/数据流的 TPL 数据流
【发布时间】:2018-09-28 11:31:09
【问题描述】:

假设您有一个配置了并行度的 TransformBlock,并希望通过该块传输数据。只有当管道可以实际开始处理输入数据时,才应创建输入数据。 (并且应该在它离开管道的那一刻被释放。)

我能做到这一点吗?如果是的话怎么办?

基本上我想要一个用作迭代器的数据源。 像这样:

public IEnumerable<Guid> GetSourceData()
{
    //In reality -> this should also be an async task -> but yield return does not work in combination with async/await ...
    Func<ICollection<Guid>> GetNextBatch = () => Enumerable.Repeat(100).Select(x => Guid.NewGuid()).ToArray();

    while (true)
    {
        var batch = GetNextBatch();
        if (batch == null || !batch.Any()) break;
        foreach (var guid in batch)
            yield return guid;
    }
}

这将导致内存中有 +- 100 条记录。好的:如果您附加到此数据源的块将它们在内存中保留一段时间,但您有机会仅获取数据的子集(/流),则更多。


一些背景信息:

我打算将此与 azure cosmos db 结合使用,其中源可以是集合中的所有对象,也可以是更改提要。不用说,我不希望所有这些对象都存储在内存中。所以这是行不通的:

using System.Threading.Tasks.Dataflow;

public async Task ExampleTask()
{
    Func<Guid, object> TheActualAction = text => text.ToString();

    var config = new ExecutionDataflowBlockOptions
    {
        BoundedCapacity = 5,
        MaxDegreeOfParallelism = 15
    };
    var throtteler = new TransformBlock<Guid, object>(TheActualAction, config);
    var output = new BufferBlock<object>();
    throtteler.LinkTo(output);

    throtteler.Post(Guid.NewGuid());
    throtteler.Post(Guid.NewGuid());
    throtteler.Post(Guid.NewGuid());
    throtteler.Post(Guid.NewGuid());
    //...
    throtteler.Complete();

    await throtteler.Completion;
}

上面的例子并不好,因为我添加了所有项目而不知道它们是否真的被转换块“使用”。另外,我并不真正关心输出缓冲区。我知道我需要将它发送到某个地方,以便等待完成,但在那之后我没有使用缓冲区。所以它应该忘记它所得到的一切......

【问题讨论】:

  • 这个可以工作,事实上它应该是这样工作的。如果目标缓冲区已满,您需要使用await SendAsync() 而不是Post 来阻止源
  • 或者您可以将迭代器函数包装在 TransformManyBlock 中。向它发送一条消息并让它循环生成消息
  • 如果您不关心输出,请使用 ActionBlock 而不是 TransformBlock
  • 简而言之,我建议您解释一下您真正想做的事情。数据流库已经提供了您所要求的。不过,有很多方法可以组合这些东西——await target.SendAsync() 的循环是发送带背压消息的最简单方法,但它从哪里获取数据呢?它会在等待时阻塞还是可以调用异步方法?这将如何在迭代器中工作?该循环也不会成为管道的一部分。然后将循环包装在 TransformManyBlock 中?
  • 我不想像这样使用批处理,目标是在 azure cosmos db 中获得一致但有限的​​ RU 使用。一个批次可能从 1000RU 开始,然后在批次结束时下降到可能 50RU,并且在开始下一批时回到 1000RU。最后,我想为给定进程获得(例如)+- 700 RU 的平均值。我认为我可以通过限制管道中 RU 最密集部分的并发任务数量来实现这一点。所有其他方块都可以随心所欲地运行。

标签: c# tpl-dataflow


【解决方案1】:

Post() 将返回 false 如果目标已满且没有阻塞。虽然这个可以用在忙等待循环中,但是很浪费。另一方面,SendAsync() 将等待目标已满:

public async Task ExampleTask()
{
    var config = new ExecutionDataflowBlockOptions
    {
        BoundedCapacity = 50,
        MaxDegreeOfParallelism = 15
    };
    var block= new ActionBlock<Guid, object>(TheActualAction, config);

    while(//some condition//)
    { 
        var data=await GetDataFromCosmosDB();
        await block.SendAsync(data);
        //Wait a bit if we want to use polling
        await Task.Delay(...);
    }

    block.Complete();
    await block.Completion;
}

【讨论】:

    【解决方案2】:

    您似乎希望以定义的并行度处理数据 (MaxDegreeOfParallelism = 15)。对于这样一个简单的需求,使用 TPL 数据流非常笨重。

    有一个非常简单而强大的模式可以解决您的问题。这是一个并行异步 foreach 循环,如下所述:https://blogs.msdn.microsoft.com/pfxteam/2012/03/05/implementing-a-simple-foreachasync-part-2/

    public static Task ForEachAsync<T>(this IEnumerable<T> source, int dop, Func<T, Task> body) 
    { 
        return Task.WhenAll( 
            from partition in Partitioner.Create(source).GetPartitions(dop) 
            select Task.Run(async delegate { 
                using (partition) 
                    while (partition.MoveNext()) 
                        await body(partition.Current); 
            })); 
    }
    

    你可以这样写:

    var dataSource = ...; //some sequence
    dataSource.ForEachAsync(15, async item => await ProcessItem(item));
    

    很简单。

    您可以使用SemaphoreSlim 动态降低 DOP。信号量充当只允许 N 个并发线程/任务进入的门。N 可以动态更改。

    因此,您将使用 ForEachAsync 作为基本主力,然后在顶部添加额外的限制和节流。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-01-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多