优化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
相关产品推荐
相关产品推荐

