【发布时间】: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