C# Rx如何处理背压:实现分页查询与Web API并发调用限制
Rx.NET背压场景解决方案
Rx.NET没有内置官方的背压处理机制,你需要的逐页拉取+单页并发限制的需求,可以通过调整数据流生成逻辑结合Rx自带的并发控制算子实现。
原有代码核心问题
- 原有
FetchRecords的实现逻辑是订阅后就会循环拉取所有分页,完全不感知下游处理速度,是导致所有数据提前被拉取的核心原因 Do(async x => await repo.Save(x))属于异步空写法,Rx不会等待该异步操作完成,Save操作的执行完全脱离了数据流生命周期,可能导致数据丢失或流程提前结束- 同时使用
Merge(1)和SemaphoreSlim做并发控制逻辑冲突,无法达到预期的3并发效果
修改后实现代码
using System; using System.Collections.Generic; using System.Linq; using System.Reactive.Linq; using System.Reactive.Threading.Tasks; using System.Threading.Tasks; using Castle.Core.Internal; using Xunit; using Xunit.Abstractions; namespace ProductValidation.CLI.Tests.Services { public class Example { private readonly ITestOutputHelper output; public Example(ITestOutputHelper output) { this.output = output; } [Fact] public async Task RunsObservableToCompletion() { var repo = new Repository(output); var client = new ServiceClient(output); // 逐页生成数据流,当前页处理完才拉取下一页 var results = Observable.Generate( 1, // 初始页码 _ => true, // 循环条件,后续通过TakeWhile判断是否还有数据 page => page + 1, // 下一页页码 page => repo.FetchPage(page).ToObservable(), // 拉取当前页 RxApp.TaskpoolScheduler) .Concat() // 确保按页顺序处理 .TakeWhile(pageProducts => !pageProducts.IsNullOrEmpty()) // 无数据时结束流 .SelectMany(pageProducts => // 单页内的记录并发调用API,限制3并发 pageProducts.ToObservable() .Select(x => client.FetchMoreInformation(x).ToObservable()) .Merge(3) ) // 等待存库完成再处理下一条 .SelectMany(result => repo.Save(result).ToObservable()); await results.LastOrDefaultAsync(); } } public class Repository { private readonly ITestOutputHelper output; public Repository(ITestOutputHelper output) { this.output = output; } public async Task<IEnumerable<int>> FetchPage(int page) { // 模拟分页查询延迟 await Task.Delay(500); output.WriteLine("Fetching page {0}", page); if (page >= 4) return Enumerable.Empty<int>(); return Enumerable.Range(1, 3).Select(_ => page); } public async Task Save(string id) { await Task.Delay(50); //模拟存库延迟 } } public class ServiceClient { private readonly ITestOutputHelper output; public ServiceClient(ITestOutputHelper output) { this.output = output; } public async Task<string> FetchMoreInformation(int id) { output.WriteLine("Calling the web client for {0}", id); await Task.Delay(1000); //模拟API调用延迟 return id.ToString(); } } }
关键逻辑说明
- 分页逻辑改为按需拉取:仅当当前页所有记录的API调用、存库操作全部完成后,才会触发下一页的拉取,完全匹配要求的执行顺序,不会提前拉取所有数据
- 直接使用
Merge(3)控制API请求并发数,无需额外引入信号量,实现更简洁且符合Rx编程范式 - 用
SelectMany等待存库操作完成,确保所有操作都纳入Rx的数据流生命周期管控,不会出现流程提前结束的问题 - 移除了不必要的调度器切换,避免出现调度混乱导致的逻辑异常
内容的提问来源于stack exchange,提问作者Roger
相关产品推荐
相关产品推荐

