Ruby使用concurrent-ruby异步调用方法无响应问题排查与优化建议
问题分析与解决方案
核心问题原因
你遇到的无响应问题,本质是主程序提前退出导致异步任务被强制终止:
- 当你调用
Processor.new.async.consume(row)时,concurrent-ruby会把任务放到后台线程池执行,但主程序在遍历完5万行CSV后,直接执行ensure块的输出就退出了。 - Ruby主进程结束时,所有后台线程会被强制销毁,异步任务根本没机会执行,所以看不到任何输出。
基础修复方案:等待异步任务全部完成
修改代码,收集所有异步任务返回的Future对象,强制主程序等待所有任务执行完毕再退出:
require 'csv' require "json" require 'byebug' require 'concurrent' class Processor include Concurrent::Async def consume(row, index) begin # 替换为实际的7个顺序API调用逻辑 puts "Processing row #{index} - #{row['某字段名']}" # 用实际字段代替,方便定位 sleep 1 if index % 2 sleep 5 if index % 3 puts "Completed row #{index}" rescue StandardError => e puts "Failed row #{index}: #{e.message}" # 这里可以加日志记录或错误上报逻辑 end end end class ImporterAsync def self.perform(master_file) futures = [] begin CSV.foreach(master_file, headers: true).with_index(1) do |row, index| # 收集每个异步任务的Future实例 future = Processor.new.async.consume(row, index) futures << future end # 等待所有异步任务执行完成 puts "All tasks submitted, waiting for completion..." Concurrent::Future.wait_all(futures) rescue StandardError => e puts "Main script failed: #{e.message}" ensure puts 'All tasks finished, closing' end end end ImporterAsync.perform('master_file.csv')
关键优化建议
针对5万行数据+多API调用的场景,只解决异步执行问题还不够,还需要考虑稳定性和效率:
1. 控制并发数,避免API限流
第三方API通常有速率限制,一次性启动5万并发会直接触发限流甚至封禁。可以自定义线程池控制并发数:
# 自定义线程池:最小5个线程,最大20个并发,任务队列最多存100个 executor = Concurrent::ThreadPoolExecutor.new( min_threads: 5, max_threads: 20, max_queue: 100 ) # 给Processor指定这个线程池 class Processor include Concurrent::Async(executor: executor) # ... 原有代码 end
根据API的实际速率限制调整max_threads数值。
2. 增加API重试机制
网络波动、API临时故障是常态,给API调用加指数退避重试:
def consume(row, index) max_retries = 3 retry_delay = 2 begin # 第1个API调用 # 第2个API调用 # ... 直到第7个 rescue Net::OpenTimeout, Net::ReadTimeout, Faraday::ServerError => e max_retries -= 1 if max_retries > 0 puts "Retrying row #{index} (remaining: #{max_retries})" sleep retry_delay retry_delay *= 2 # 指数退避,每次等待时间翻倍 else puts "Failed after retries: row #{index}, error: #{e.message}" end end end
3. 替换puts为日志系统
用Ruby标准库Logger记录详细日志,方便排查问题:
require 'logger' class Processor @@logger = Logger.new('import_errors.log', 'daily') # 按天滚动日志 def consume(row, index) begin @@logger.info("Start processing row #{index}") # API调用逻辑 @@logger.info("Finish processing row #{index}") rescue StandardError => e @@logger.error("Row #{index} failed: #{e.message}\n#{e.backtrace.join("\n")}") end end end
4. 进度跟踪
添加计数器实时查看处理进度:
class ImporterAsync def self.perform(master_file) futures = [] processed = Concurrent::AtomicFixnum.new(0) total_rows = CSV.read(master_file, headers: true).size # 先获取总行数 CSV.foreach(master_file, headers: true).with_index(1) do |row, index| future = Processor.new.async.consume(row, index) future.add_callback do count = processed.increment puts "Progress: #{count}/#{total_rows} (#{(count/total_rows.to_f*100).round(2)}%)" end futures << future end Concurrent::Future.wait_all(futures) end end
内容的提问来源于stack exchange,提问作者Ajith
相关产品推荐
相关产品推荐

