【发布时间】:2021-09-06 07:33:45
【问题描述】:
我正在编写一个批处理管道,它每 Y 秒处理 X 个未完成的操作。感觉System.Reactive 很适合这个,但我无法让订阅者并行执行。我的代码如下所示:
var subject = new Subject<int>();
var concurrentCount = 0;
using var reader = subject
.Buffer(TimeSpan.FromSeconds(1), 100)
.Subscribe(list =>
{
var c = Interlocked.Increment(ref concurrentCount);
if (c > 1) Console.WriteLine("Executing {0} simultaneous batches", c); // This never gets printed, because Subscribe is only ever called on a single thread.
Interlocked.Decrement(ref concurrentCount);
});
Parallel.For(0, 1_000_000, i =>
{
subject.OnNext(i);
});
subject.OnCompleted();
有没有一种优雅的方式以并发方式从这个缓冲的Subject 中读取?
【问题讨论】:
标签: c# .net system.reactive