.NET Core中Azure服务总线触发函数调用第三方API限流异常排查
问题分析与解决方案
你的限流逻辑存在几个核心问题,导致最后几条消息无法正常处理:
内存计数不可靠
你用的requestCount是内存变量,Azure函数是无状态的,哪怕设置了最大并发为1,函数实例重启、冷启动或者回收都会导致计数丢失;另外,当你在else块里延迟60秒时,Service Bus的消息锁可能已经过期(默认锁时长通常是30秒),消息会被重新放回队列,造成重复处理或计数混乱。计数逻辑时机错误
当前代码里,只有当requestCount <45时才调用API,但没有在调用成功后递增计数;而且进入else块后,重置计数并调用API,但这条调用的计数没有被记录,导致后续消息的计数逻辑彻底混乱。延迟处理导致消息锁超时
使用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
相关产品推荐
相关产品推荐

