You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.05 15:15:03