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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 21:10:21