如何解决多Sidekiq任务中的API速率限制问题?
缓解Sidekiq批量任务API速率限制的实用方案
1. 批量触发时做限流控制
不要一次性把所有任务丢进Sidekiq队列,而是分批次、带间隔地触发:
- 手动拆分批量任务,比如每处理50个任务就暂停10秒,避免瞬间占满API配额:
def perform(batch_items) batch_items.each_slice(50) do |slice| slice.each { |item| ExternalApiWorker.perform_async(item) } sleep(10) # 批次间隔,根据API配额调整 end end
- 借助
sidekiq-rate-limitergem,给调用外部API的任务队列设置全局速率限制,比如限制每分钟最多执行60次:
# 在ExternalApiWorker里配置 sidekiq_options rate_limit: { limit: 60, period: 60 }
2. 利用API速率限制响应动态调整重试
大多数API会在429响应头返回配额重置时间,直接根据这个时间延迟任务,避免盲目重试:
class ExternalApiWorker include Sidekiq::Worker def perform(api_params) response = Faraday.post('https://external-api.com/endpoint', api_params) handle_response(response, api_params) rescue Faraday::Error::ClientError => e if e.response&.status == 429 reset_epoch = e.response.headers['X-RateLimit-Reset'].to_i delay_seconds = [reset_epoch - Time.now.to_i + 1, 10].max # 至少延迟10秒 self.class.perform_in(delay_seconds.seconds, api_params) else # 其他错误按正常逻辑处理 raise e end end private def handle_response(response, api_params) # 正常业务逻辑处理 end end
3. 自定义重试策略替代默认指数退避
放弃Sidekiq默认的指数退避,改成更贴合API配额的重试间隔,同时限制最大重试次数:
class ExternalApiWorker include Sidekiq::Worker sidekiq_options retry: 3, retry_in: ->(count) { # 根据重试次数逐步增加间隔,同时参考API重置时间 base_delay = [10, 30, 60][count] || 120 # 如果能拿到API重置时间,优先用这个值 reset_time = get_api_reset_time reset_time ? (reset_time - Time.now.to_i + 1) : base_delay } private def get_api_reset_time # 调用API的HEAD请求获取当前配额重置时间,或者从缓存读取 response = Faraday.head('https://external-api.com/endpoint') response.headers['X-RateLimit-Reset'].to_i if response.status == 200 rescue nil end end
4. 单独队列+低并发限制
把调用外部API的任务放到专属队列,在sidekiq.yml里限制该队列的worker并发数,确保同一时间只有少量任务在调用API:
:concurrency: 10 :queues: - [critical, 8] - [external_api, 2] # 限制该队列最多2个并发worker
5. 合并请求减少API调用次数
如果外部API支持批量接口,把多个单任务的请求合并成一次批量请求,大幅降低调用频率:
# 原单任务worker class ExternalApiWorker include Sidekiq::Worker def perform(item_id) item = Item.find(item_id) ExternalApi.call(item.data) end end # 改成批量worker class BatchExternalApiWorker include Sidekiq::Worker def perform(item_ids) items = Item.where(id: item_ids) batch_data = items.map(&:data) ExternalApi.batch_call(batch_data) # 调用批量接口 end end # 批量触发时调用批量worker def perform(batch_items) batch_items.each_slice(100) do |slice| BatchExternalApiWorker.perform_async(slice.map(&:id)) sleep(5) end end
内容的提问来源于stack exchange,提问作者paulo vilarinho
相关产品推荐
相关产品推荐

