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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 06:10:18