System.Reactive中如何实现缓冲Subject的订阅者并行执行?
System.Reactive 缓冲批处理并行执行方案
问题本质
Rx 默认遵循观察者通知序列化约定,同一订阅者的 OnNext 回调永远是串行调用的,不会出现并发执行的情况,这是原代码中 concurrentCount 永远无法大于 1 的根本原因。
优雅实现方案
推荐使用 Rx 原生的 FromAsync + 指定并行度的 Merge 操作符实现,既符合 Rx 编程范式,还能灵活控制最大并行批数:
var subject = new Subject<int>(); int concurrentCount = 0; // 可根据业务需求调整最大并行批处理数量 int maxConcurrentBatches = 4; using var reader = subject .Buffer(TimeSpan.FromSeconds(1), 100) // 将每个批次的处理逻辑封装为冷Observable,仅被订阅时才执行 .Select(buffer => Observable.FromAsync(async () => { var c = Interlocked.Increment(ref concurrentCount); if (c > 1) Console.WriteLine("Executing {0} simultaneous batches", c); // 替换为实际的批处理业务逻辑,示例用Delay模拟处理耗时 await Task.Delay(100); Interlocked.Decrement(ref concurrentCount); })) // Merge操作符控制同时执行的批处理任务数量 .Merge(maxConcurrentBatches) .Subscribe(); Parallel.For(0, 1_000_000, i => { subject.OnNext(i); }); subject.OnCompleted();
简化备选方案
如果不需要控制并行度,也可以直接在 Subscribe 回调中将处理逻辑丢到线程池执行,让 OnNext 立刻返回处理下一个批次:
var subject = new Subject<int>(); var concurrentCount = 0; using var reader = subject .Buffer(TimeSpan.FromSeconds(1), 100) .Subscribe(list => { // 丢弃Task返回值,不等待执行完成直接返回 _ = Task.Run(() => { var c = Interlocked.Increment(ref concurrentCount); if (c > 1) Console.WriteLine("Executing {0} simultaneous batches", c); Interlocked.Decrement(ref concurrentCount); }); }); Parallel.For(0, 1_000_000, i => { subject.OnNext(i); }); subject.OnCompleted();
该方案实现简单,但无法控制并行上限,批次积压时可能导致线程数过载,仅适合低负载场景使用。
内容的提问来源于stack exchange,提问作者Ceilingfish
相关产品推荐
相关产品推荐

