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

优化ConfirmPendingAccountsJob批量上传性能的技术问询

批量账号上传性能优化需求

当前正在重构账号批量上传的Job性能,现有账号上传耗时过长,已对数据库负载和性能造成影响。现需确认ConfirmPendingAccountsJob能否高效处理至少50,000条账号,同时需要实现以下优化:

  • 将批量账号传递给下游方法(重点是team.add_accounts和licence.bulk_enrol_accounts)以提升处理效率
  • 实现统一的分批处理逻辑,让Job可接收任意数量的账号,单个Job仅处理可控数量的账号
  • 支持配置批量大小和批次间隔参数,用于保护数据库

现有调用逻辑

def confirm_pending_account(account_upload)
  ConfirmPendingAccountsJob.perform_later(account_upload: account_upload)
end

现有ConfirmPendingAccountsJob实现

class ConfirmPendingAccountsJob < BaseJob
  queue_as :low_priority

  attr_reader :accounts, :account_upload

  def perform(account_upload:)
    @account_upload = account_upload
    pending_accounts = account_upload.pending_accounts.not_failed

    @accounts = {
      previously_existing: [],
      newly_created: []
    }

    pending_accounts.each do |pending_account|
      # TODO: Find duplicates in bulk
      duplicate_account_in_same_organisation = pending_account.find_duplicate
      if duplicate_account_in_same_organisation
        accounts[:previously_existing] << [pending_account, duplicate_account_in_same_organisation]
        next
      end

      # Verified is set before this job is run; the next line is only needed in order to run this
      # job in isolation, from the command line, otherwise all account creation will be skipped
      # pending_account.account_verified!
      next unless pending_account.verified? # Failure already recorded by VerifyPendingAccountJob

      created_account = pending_account.create_account_with_errors
      if created_account.errors.any?
        pending_account.record_failure!(serialized_errors: created_account.errors.to_hash)
      else
        accounts[:newly_created] << [pending_account, created_account]
      end
    end

    finalise_accounts_creation
  end

  private

  def finalise_accounts_creation
    move_into_teams_and_licences

    accounts[:previously_existing].each do |bundle|
      update_pending_account(bundle.pending_account, bundle.account, status: :updated)
    end

    accounts[:newly_created].each do |pending_account, account|
      update_pending_account(pending_account, account, status: :created)
      account.confirm
    end
  end

  def move_into_teams_and_licences
    Team.find(account_upload.team_ids).each do |team|
      team.add_accounts(all_accounts)
    end

    Licensing::Public::API.find_licence(account_upload.licence_ids).each do |licence|
      licence.bulk_enrol_accounts(all_accounts)
    end

    EventPublisher.account_licences_changed(all_accounts)
  end

  def update_pending_account(pending_account, account, status:)
    if status == :created
      pending_account.create_account!
    else
      pending_account.update_account!
    end
    # This establishes the link from an account to a pending account
    pending_account.update!(account: account)
  end

  def all_accounts
    newly_created_accounts + previously_existing_accounts
  end

  def previously_existing_accounts
    @accounts[:previously_existing].map{|bundle| bundle.account }
  end

  def newly_created_accounts
    @accounts[:newly_created].map{|_,account| account }
  end
end

优化方案

1. 实现分批处理逻辑

修改Job支持分批处理,通过环境变量配置批次大小和间隔,避免单Job处理过大数据量:

class ConfirmPendingAccountsJob < BaseJob
  queue_as :low_priority

  # 可配置参数,根据数据库性能调整默认值
  BATCH_SIZE = ENV.fetch('ACCOUNT_CONFIRM_BATCH_SIZE', 500).to_i
  BATCH_DELAY = ENV.fetch('ACCOUNT_CONFIRM_BATCH_DELAY', 10).to_i # 单位:秒

  attr_reader :accounts, :account_upload, :offset

  def perform(account_upload:, offset: 0)
    @account_upload = account_upload
    @offset = offset

    # 分批获取待处理账号
    pending_accounts = account_upload.pending_accounts.not_failed.offset(offset).limit(BATCH_SIZE)
    return if pending_accounts.empty? # 无待处理账号时终止

    @accounts = {
      previously_existing: [],
      newly_created: []
    }

    process_pending_accounts(pending_accounts)
    finalise_accounts_creation

    # 调度下一批次Job,间隔指定时间
    self.class.perform_in(BATCH_DELAY.seconds, account_upload: account_upload, offset: offset + BATCH_SIZE)
  end

  private

  # 提取原有账号处理逻辑为单独方法
  def process_pending_accounts(pending_accounts)
    pending_accounts.each do |pending_account|
      duplicate_account_in_same_organisation = pending_account.find_duplicate
      if duplicate_account_in_same_organisation
        accounts[:previously_existing] << [pending_account, duplicate_account_in_same_organisation]
        next
      end

      next unless pending_account.verified?

      created_account = pending_account.create_account_with_errors
      if created_account.errors.any?
        pending_account.record_failure!(serialized_errors: created_account.errors.to_hash)
      else
        accounts[:newly_created] << [pending_account, created_account]
      end
    end
  end

  # 优化下游方法调用:传递当前批次账号而非全部
  def move_into_teams_and_licences
    current_batch_accounts = all_accounts
    return if current_batch_accounts.empty?

    Team.find(account_upload.team_ids).each do |team|
      team.add_accounts(current_batch_accounts)
    end

    Licensing::Public::API.find_licence(account_upload.licence_ids).each do |licence|
      licence.bulk_enrol_accounts(current_batch_accounts)
    end

    EventPublisher.account_licences_changed(current_batch_accounts)
  end

  # 其余原有方法保持不变...
end

2. 落地批量查重逻辑(解决TODO项)

原代码单账号查重会导致N+1查询,替换为批量查询减少数据库请求:

def process_pending_accounts(pending_accounts)
  # 提取待处理账号的查重标识(示例:假设用邮箱作为查重字段)
  emails = pending_accounts.map(&:email)
  # 批量查询组织内的重复账号,按邮箱建立索引
  duplicate_accounts = Account.where(organisation_id: account_upload.organisation_id)
                               .where(email: emails)
                               .index_by(&:email)

  pending_accounts.each do |pending_account|
    duplicate_account = duplicate_accounts[pending_account.email]
    if duplicate_account
      accounts[:previously_existing] << [pending_account, duplicate_account]
      next
    end

    # 后续账号创建逻辑保持不变...
  end
end

3. 配置参数说明

  • ACCOUNT_CONFIRM_BATCH_SIZE:单个Job处理的账号数量,默认500,可根据数据库CPU、IO负载调整
  • ACCOUNT_CONFIRM_BATCH_DELAY:批次之间的间隔时间,默认10秒,用于给数据库留出资源恢复窗口

4. 下游方法适配

确保team.add_accounts和licence.bulk_enrol_accounts本身支持批量操作(比如使用insert_all创建关联,而非循环调用单条创建方法),否则分批的性能优势会被下游的低效逻辑抵消。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 13:20:28