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

C# Rx如何处理背压:实现分页查询与Web API并发调用限制

Rx.NET背压场景解决方案

Rx.NET没有内置官方的背压处理机制,你需要的逐页拉取+单页并发限制的需求,可以通过调整数据流生成逻辑结合Rx自带的并发控制算子实现。

原有代码核心问题

  1. 原有FetchRecords的实现逻辑是订阅后就会循环拉取所有分页,完全不感知下游处理速度,是导致所有数据提前被拉取的核心原因
  2. Do(async x => await repo.Save(x))属于异步空写法,Rx不会等待该异步操作完成,Save操作的执行完全脱离了数据流生命周期,可能导致数据丢失或流程提前结束
  3. 同时使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 00:06:06