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

.NET Core中Azure服务总线触发函数调用第三方API限流异常排查

问题分析与解决方案

你的限流逻辑存在几个核心问题,导致最后几条消息无法正常处理:

  1. 内存计数不可靠
    你用的requestCount是内存变量,Azure函数是无状态的,哪怕设置了最大并发为1,函数实例重启、冷启动或者回收都会导致计数丢失;另外,当你在else块里延迟60秒时,Service Bus的消息锁可能已经过期(默认锁时长通常是30秒),消息会被重新放回队列,造成重复处理或计数混乱。

  2. 计数逻辑时机错误
    当前代码里,只有当requestCount <45时才调用API,但没有在调用成功后递增计数;而且进入else块后,重置计数并调用API,但这条调用的计数没有被记录,导致后续消息的计数逻辑彻底混乱。

  3. 延迟处理导致消息锁超时
    使用Thread.Sleep或Task.Delay会让函数线程阻塞/等待,这段时间内Service Bus的消息锁如果没有续期,消息会被视为未处理,重新回到队列,最终导致部分消息重复投递或遗留。


修复后的实现思路

1. 改用分布式计数器(比如Azure Redis Cache)

用Redis的原子操作(INCR、EXPIRE)来维护每分钟的请求计数,确保计数在函数实例之间共享且不会丢失:

// 生成当前分钟的唯一键,保证每分钟计数独立
var currentMinuteKey = $"api_requests_{DateTime.UtcNow:yyyyMMddHHmm}";
var redis = ConnectionMultiplexer.Connect("<你的Redis连接字符串>").GetDatabase();

// 原子递增计数,同时设置键1分钟后自动过期
var requestCount = await redis.StringIncrementAsync(currentMinuteKey);
await redis.KeyExpireAsync(currentMinuteKey, TimeSpan.FromMinutes(1));

if (requestCount <= 45)
{
    var response = await MakeRequest(json);
    if (!response.IsSuccessStatusCode)
    {
        // API调用失败,递减计数(如果失败请求不计入限流配额)
        await redis.StringDecrementAsync(currentMinuteKey);
        // 抛出异常让Service Bus重试该消息
        throw new HttpRequestException($"API调用失败,状态码:{response.StatusCode}");
    }
}
else
{
    // 计算到下一分钟的剩余等待时间,不用硬等60秒
    var now = DateTime.UtcNow;
    var nextMinute = now.AddMinutes(1).AddSeconds(-now.Second).AddMilliseconds(-now.Millisecond);
    var delay = nextMinute - now;
    await Task.Delay(delay);

    // 进入下一分钟后重新计数
    var newMinuteKey = $"api_requests_{nextMinute:yyyyMMddHHmm}";
    await redis.StringIncrementAsync(newMinuteKey);
    await redis.KeyExpireAsync(newMinuteKey, TimeSpan.FromMinutes(1));

    var response = await MakeRequest(json);
    if (!response.IsSuccessStatusCode)
    {
        await redis.StringDecrementAsync(newMinuteKey);
        throw new HttpRequestException($"API调用失败,状态码:{response.StatusCode}");
    }
}

2. 配置Service Bus消息锁时长

在函数的host.json里延长消息锁时长并开启自动续锁,避免延迟处理时锁过期:

{
  "version": "2.0",
  "extensions": {
    "serviceBus": {
      "messageHandlerOptions": {
        "autoComplete": false, // 关闭自动完成,手动控制消息状态
        "maxConcurrentCalls": 1,
        "maxAutoRenewDuration": "00:05:00" // 自动续锁最长5分钟,覆盖延迟时间
      }
    }
  }
}

注意在函数处理完成后,要手动调用message.CompleteAsync()来标记消息已处理;如果处理失败,调用message.AbandonAsync()让消息重新进入队列重试。

3. 用Polly实现限流(更简洁可靠)

结合Polly的RateLimiter策略,配合分布式缓存实现全局限流,代码更简洁且维护性更高:

// 初始化基于Redis的分布式限流策略
var redis = ConnectionMultiplexer.Connect("<你的Redis连接字符串>").GetDatabase();
var rateLimiter = PartitionedRateLimitStrategy.Create<string>(
    partitionKey => RateLimitStrategy.Create(
        SlidingWindowRateLimiterOptions.Create(
            limit: 45,
            window: TimeSpan.FromMinutes(1),
            queueProcessingOrder: QueueProcessingOrder.OldestFirst
        ),
        new RedisRateLimiterStore(redis, partitionKey)
    )
);

// 执行限流调用
await rateLimiter.ExecuteAsync(
    partitionKey: "api_requests",
    async token =>
    {
        var response = await MakeRequest(json);
        if (!response.IsSuccessStatusCode)
        {
            throw new HttpRequestException($"API调用失败,状态码:{response.StatusCode}");
        }
    }
);

关键注意点

  • 永远不要用内存变量维护跨函数调用的状态,必须用分布式存储(Redis、Azure Table Storage等)。
  • 处理API失败的情况,避免失败的请求占用限流配额。
  • 确保Service Bus消息锁的时长足够覆盖延迟处理的时间,开启自动续锁机制。

内容的提问来源于stack exchange,提问作者Tech_Explorer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 08:03:21