【问题标题】:Big task or multiple small tasks with Sidekiq使用 Sidekiq 完成大任务或多个小任务
【发布时间】:2015-08-25 21:52:20
【问题描述】:

我正在编写一个工人来将很多用户添加到一个组中。我想知道是否最好运行一个拥有所有用户的大型任务,或者像 100 个用户一样批量运行,或者每个任务一个接一个地运行。

这里是我的代码

class AddUsersToGroupWorker
  include Sidekiq::Worker
  sidekiq_options :queue => :group_utility

  def perform(store_id, group_id, user_ids_to_add)
    begin
      store = Store.find store_id
      group = Group.find group_id
    rescue ActiveRecord::RecordNotFound => e
      Airbrake.notify e
      return
    end

    users_to_process = store.users.where(id: user_ids_to_add)
                                  .where.not(id: group.user_ids)
    group.users += users_to_process

    users_to_process.map(&:id).each do |user_to_process_id|
      UpdateLastUpdatesForUserWorker.perform_async store.id, user_to_process_id
    end
  end
end 

也许在我的方法中有这样的东西会更好:

def add_users
    users_to_process = store.users.where(id: user_ids_to_add)
                                  .where.not(id: group.user_ids)

    users_to_process.map(&:id).each do |user_to_process_id|
      AddUserToGroupWorker.perform_async group_id, user_to_process_id
      UpdateLastUpdatesForUserWorker.perform_async store.id, user_to_process_id
    end
end

但是有很多find 请求。你怎么看?

如果需要,我有一个 sidekig pro 许可证(例如批量)。

【问题讨论】:

  • 你的UpdateLastUpdatesForUserWorker 是做什么的?
  • 嘿Eugzol!它向用户发送通知,并在组中附加与用户相关的内容。
  • 再次嗨 :) 发送通知——什么通知?某种远程请求,例如发送电子邮件或推送通知?当您对用户或组说 attach 某事时,您是什么意思?
  • 是的移动通知和电子邮件。如果他在特定组中,用户可以阅读帖子。他还可以访问文件。 UpdateLastUpdatesForUserWorker 被许多其他进程使用。

标签: ruby sidekiq


【解决方案1】:

这是我的想法。

1.执行单个 SQL 查询而不是 N 个查询

这一行:group.users += users_to_process 可能会产生 N 个 SQL 查询(其中 N 是 users_to_process.count)。我假设您在用户和组之间有多对多连接(使用user_groups 连接表/模型),所以您应该使用一些Mass inserting data technique

users_to_process_ids = store.users.where(id: user_ids_to_add)
                         .where.not(id: group.user_ids)
                         .pluck(:id)
sql_values = users_to_process_ids.map{|i| "(#{i.to_i}, #{group.id.to_i}, NOW(), NOW())"}
Group.connection.execute("
  INSERT INTO groups_users (user_id, group_id, created_at, updated_at)
  VALUES #{sql_values.join(",")}
")

是的,它是原始 SQL。而且它很快

2。用户 pluck(:id) 而不是 map(&:id)

pluck 更快,因为:

  • 它将仅选择“id”列,因此从 DB 传输的数据较少
  • 更重要的是,它不会为每个原始数据创建 ActiveRecord 对象

执行 SQL 很便宜。创建 Ruby 对象非常昂贵。

3。使用水平并行化而不是垂直并行化

我这里的意思是,如果你需要为十几条记录做顺序任务A -> B -> C,有两种主要的工作拆分方式:

  • 垂直分割AWorkerA(1), A(2), A(3); BWorkerB(1) 等; CWorker 完成所有 C(i) 工作;
  • 水平分割UniversalWorkerA(1)+B(1)+C(1)

使用后一种(水平)方式。

这是来自经验的陈述,而不是从某种理论的角度来看(两种方式都可行)。

为什么要这样做?

  • 当您使用垂直分段时,您可能会在将工作从一名工作人员传递给另一名工作人员时出错。喜欢such kind of errors。如果你遇到这样的错误,你会拔掉头发,因为它们不是持久的,也不容易重现。有时它们会发生,有时它们不会。是否可以编写一个代码,将工作沿着链向下传递而不会出错?是的。不过最好keep it simple
  • 假设您的服务器处于静止状态。然后突然新的工作来了。你的BC 工作人员只会浪费内存,而你的A 工作人员会完成这项工作。然后你的AC 将浪费内存,而B 正在工作。等等。如果您进行水平分割,您的资源消耗会自行解决。

将该建议应用于您的具体案例:对于初学者,不要在另一个异步任务中调用 perform_async

4.分批处理

回答您最初的问题 - 是的,分批处理。创建和管理异步任务本身会占用一些资源,因此无需创建太多。


TL;DR所以最后,您的代码可能看起来像这样:

# model code

BATCH_SIZE = 100

def add_users
  users_to_process_ids = store.users.where(id: user_ids_to_add)
                           .where.not(id: group.user_ids)
                           .pluck(:id)
  # With 100,000 users performance of this query should be acceptable
  # to make it in a synchronous fasion
  sql_values = users_to_process_ids.map{|i| "(#{i.to_i}, #{group.id.to_i}, NOW(), NOW())"}
  Group.connection.execute("
    INSERT INTO groups_users (user_id, group_id, created_at, updated_at)
    VALUES #{sql_values.join(",")}
  ")

  users_to_process_ids.each_slice(BATCH_SIZE) do |batch|
    AddUserToGroupWorker.perform_async group_id, batch
  end
end

# add_user_to_group_worker.rb

def perform(group_id, user_ids_to_add)
  group = Group.find group_id

  # Do some heavy load with a batch as a whole
  # ...
  # ...
  # If nothing here is left, call UpdateLastUpdatesForUserWorker from the model instead

  user_ids_to_add.each do |id|
    # do it synchronously – we already parallelized the job
    # by splitting it in slices in the model above
    UpdateLastUpdatesForUserWorker.new.perform store.id, user_to_process_id
  end
end

【讨论】:

  • 再次感谢@EugZol 提供这个有据可查的答案。我一定会记住并行化和所有其他技巧。你不认为从模型中调用 UpdateLastUpdatesForUserWorker worker 就够了吗?成本是工人本身而不是呼叫否?
  • 不客气。我不明白你的问题 - 你能详细说明一下吗?您可以从模型中调用UpdateLastUpdatesForUserWorker。我建议不要为每个 user_id 调用perform_async,为一堆 id 调用一次perform_async
  • not to call perform_async for every user_id, call perform_async once for a bunch of ids 我想知道。谢谢
  • 我终于做了类似Group.connection.execute("INSERT INTO groups_users (user_id, group_id) VALUES #{sql_values.join(",")}") 的操作,但无法让用户正确加入我的组。 add_user 请求做 Group.last.add_user User.last : SQL (4.0ms) INSERT INTO "groups_users" ("group_id", "user_id", "created_at", "updated_at") VALUES ($1, $2, $3, $4) RETURNING "id" [["group_id", 8986], ["user_id", 119506], ["created_at", "2015-08-25 13:15:26.421200"], ["updated_at", "2015-08-25 13:15:26.421200"]]
  • 您的请求会产生什么?即,"INSERT INTO groups_users (user_id, group_id) VALUES #{sql_values.join(",")} 的编译值是多少?
【解决方案2】:

没有灵丹妙药。这取决于您的目标和您的应用程序。要问自己的一般问题:

  • 您可以将多少用户 ID 传递给工作人员?能过100吗? 100万呢?
  • 您的工人可以工作多长时间?它应该对工作时间有任何限制吗?他们能卡住吗?

对于大型应用程序,有必要将传递的参数拆分为更小的块,以避免创建长时间运行的作业。创建大量小型作业可以让您轻松扩展 - 您可以随时添加更多工人。

另外,为工作人员定义超时类型以停止处理卡住的工作人员可能是一个好主意。

【讨论】:

  • 谢谢亚历山大。处理的最大用户数约为 100.000。工作时间没有限制。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-03-03
  • 1970-01-01
  • 2012-07-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-12
相关资源
最近更新 更多