WebAPI请求延迟与排队实现:应对下游API速率限制
请求排队与速率控制实现方案
针对下游API每秒仅允许一次调用的限制,我们可以通过异步队列+后台处理器的方式实现请求排队与间隔延迟,确保所有请求按序处理且严格遵守速率限制。以下提供F#和C#两种实现方式:
F# 实现
1. 创建速率限制队列服务
定义单例后台服务,管理请求队列并按每秒一次的速率处理请求:
open System open System.Threading open System.Threading.Tasks open System.Threading.Channels open Microsoft.Extensions.Hosting open Microsoft.Extensions.Logging type WeighScanRequest = { OrderNumber: string Weight: decimal CarrierServiceOverride: string Suffix: string CompletionSource: TaskCompletionSource<string> } type RateLimitedQueueService(logger: ILogger<RateLimitedQueueService>) = inherit BackgroundService() // 创建有界队列,避免内存溢出(可根据需求调整容量) let channel = Channel.CreateBounded<WeighScanRequest>(100) member _.EnqueueRequest(request: WeighScanRequest) = channel.Writer.WriteAsync(request).AsTask() override _.ExecuteAsync(cancellationToken) = task { logger.LogInformation("Rate-limited queue service started") while not cancellationToken.IsCancellationRequested do // 从队列中取出下一个请求 let! request = channel.Reader.ReadAsync(cancellationToken) try // 调用下游API let! result = WeighScanM2.weighScanM2( request.OrderNumber, request.Weight, request.CarrierServiceOverride, request.Suffix) request.CompletionSource.SetResult(result) with | ex -> logger.LogError(ex, "Failed to process weigh scan request for order {OrderNumber}", request.OrderNumber) request.CompletionSource.SetException(ex) // 等待1秒再处理下一个请求 do! Task.Delay(TimeSpan.FromSeconds(1.0), cancellationToken) }
2. 修改控制器代码
注入队列服务,将请求提交到队列并等待处理结果:
[<ApiController>] [<Route("[controller]")>] type WeighScanController (logger : ILogger<WeighScanController>, queueService: RateLimitedQueueService) = inherit ControllerBase() [<HttpGet>] member _.Get(orderNumber: string, weight: decimal, carrierServiceOverride: string, orderNumberSuffix: string ) = task { if String.IsNullOrEmpty(orderNumber) || weight = 0m then return "Incorrect Values Passed" else let suffix = if String.IsNullOrEmpty(orderNumberSuffix) then "" else orderNumberSuffix // 创建任务完成源,接收处理结果 let tcs = TaskCompletionSource<string>() let request = { OrderNumber = orderNumber Weight = weight CarrierServiceOverride = carrierServiceOverride Suffix = suffix CompletionSource = tcs } // 将请求加入队列 do! queueService.EnqueueRequest(request) // 等待队列处理完成并返回结果 return! tcs.Task }
3. 注册服务
在Program.fs中注册后台服务:
builder.Services.AddHostedService<RateLimitedQueueService>()
C# 实现
1. 创建速率限制队列服务
using System; using System.Threading; using System.Threading.Tasks; using System.Threading.Channels; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; public class WeighScanRequest { public string OrderNumber { get; set; } public decimal Weight { get; set; } public string CarrierServiceOverride { get; set; } public string Suffix { get; set; } public TaskCompletionSource<string> CompletionSource { get; set; } } public class RateLimitedQueueService : BackgroundService { private readonly ILogger<RateLimitedQueueService> _logger; private readonly Channel<WeighScanRequest> _channel; public RateLimitedQueueService(ILogger<RateLimitedQueueService> logger) { _logger = logger; // 创建有界队列,限制最大容量 _channel = Channel.CreateBounded<WeighScanRequest>(100); } public async Task EnqueueRequest(WeighScanRequest request) { await _channel.Writer.WriteAsync(request); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation("Rate-limited queue service started"); while (!stoppingToken.IsCancellationRequested) { var request = await _channel.Reader.ReadAsync(stoppingToken); try { var result = await WeighScanM2.weighScanM2( request.OrderNumber, request.Weight, request.CarrierServiceOverride, request.Suffix); request.CompletionSource.SetResult(result); } catch (Exception ex) { _logger.LogError(ex, "Failed to process weigh scan request for order {OrderNumber}", request.OrderNumber); request.CompletionSource.SetException(ex); } // 间隔1秒处理下一个请求 await Task.Delay(TimeSpan.FromSeconds(1), stoppingToken); } } }
2. 修改控制器代码
using Microsoft.AspNetCore.Mvc; using Microsoft.Extensions.Logging; using System.Threading.Tasks; [ApiController] [Route("[controller]")] public class WeighScanController : ControllerBase { private readonly ILogger<WeighScanController> _logger; private readonly RateLimitedQueueService _queueService; public WeighScanController(ILogger<WeighScanController> logger, RateLimitedQueueService queueService) { _logger = logger; _queueService = queueService; } [HttpGet] public async Task<IActionResult> Get(string orderNumber, decimal weight, string carrierServiceOverride, string orderNumberSuffix) { if (string.IsNullOrEmpty(orderNumber) || weight == 0) { return Ok("Incorrect Values Passed"); } var suffix = string.IsNullOrEmpty(orderNumberSuffix) ? "" : orderNumberSuffix; var tcs = new TaskCompletionSource<string>(); var request = new WeighScanRequest { OrderNumber = orderNumber, Weight = weight, CarrierServiceOverride = carrierServiceOverride, Suffix = suffix, CompletionSource = tcs }; await _queueService.EnqueueRequest(request); var result = await tcs.Task; return Ok(result); } }
3. 注册服务
在Program.cs中注册后台服务:
builder.Services.AddHostedService<RateLimitedQueueService>();
关键说明
- 使用
Channel实现异步队列,确保请求按序处理,同时通过有界队列支持背压,避免内存耗尽。 - 后台服务
RateLimitedQueueService负责控制调用速率,每秒处理一个请求。 - 控制器将请求提交到队列后,通过
TaskCompletionSource异步等待处理结果,不阻塞请求线程。
内容的提问来源于stack exchange,提问作者Martin Thompson
相关产品推荐
相关产品推荐

