【问题标题】:Celery parallel distributed task with multiprocessing具有多处理功能的 Celery 并行分布式任务
【发布时间】:2014-07-17 23:01:54
【问题描述】:

我有一个 CPU 密集型 Celery 任务。我想使用跨多个 EC2 实例的所有处理能力(核心)来更快地完成这项工作(具有多处理功能的 celery 并行分布式任务 - 我认为

threadingmultiprocessingdistributed Computingdistributed parallel processing这些术语都是我的术语试图更好地理解。

示例任务:

  @app.task
  for item in list_of_millions_of_ids:
      id = item # do some long complicated equation here very CPU heavy!!!!!!! 
      database.objects(newid=id).save()

使用上面的代码(如果可能的话,还有一个例子)之前人们会如何使用 Celery 分配这个任务,方法是允许使用所有的计算云中所有可用机器的 CPU 能力?

【问题讨论】:

标签: python django multithreading multiprocessing celery


【解决方案1】:

你的目标是:

  1. 将您的工作分配给多台机器(分布式 计算/分布式并行处理)
  2. 将给定机器上的工作分配给所有 CPU (多处理/线程)

Celery 可以很容易地为您做到这两点。首先要了解的是,每个 celery worker 都是 configured by default 来运行与系统上可用的 CPU 内核一样多的任务:

Concurrency 是用于处理的 prefork 工作进程的数量 你的任务同时进行,当所有这些都忙于做新的工作时 任务必须等待其中一项任务完成才能完成 进行处理。

默认并发数是该机器上的 CPU 数量 (包括核心),您可以使用 -c 选项指定自定义编号。 没有推荐值,因为最佳数量取决于 因素的数量,但如果您的任务主要受 I/O 限制,那么您可以 尝试增加它,实验表明增加超过 CPU 数量的两倍很少有效,并且可能会降级 性能。

这意味着每个单独的任务无需担心使用多处理/线程来利用多个 CPU/内核。相反,celery 会同时运行足够多的任务来使用每个可用的 CPU。

除此之外,下一步是创建一个任务来处理您的list_of_millions_of_ids 的某些子集。您在这里有几个选项 - 一个是让每个任务处理一个 ID,因此您运行 N 个任务,其中N == len(list_of_millions_of_ids)。这将保证工作在所有任务中均匀分配,因为永远不会有一个工作人员提前完成并只是等待的情况;如果它需要工作,它可以从队列中拉出一个 id。您可以使用 celery group 执行此操作(如 John Doe 所述)。

tasks.py:

@app.task
def process_ids(item):
    id = item #long complicated equation here
    database.objects(newid=id).save()

并执行任务:

from celery import group
from tasks import process_id

jobs = group(process_ids(item) for item in list_of_millions_of_ids)
result = jobs.apply_async()

另一种选择是将列表分成更小的部分并将这些部分分发给您的工人。这种方法存在浪费一些周期的风险,因为您最终可能会导致一些工人在等待,而其他工人仍在工作。然而,celery documentation notes 认为这种担忧往往是没有根据的:

有些人可能会担心将任务分块会导致 并行性,但这对于繁忙的集群和在 练习,因为您避免了消息传递的开销 大大提高性能。

因此,由于减少了消息传递开销,您可能会发现将列表分块并将块分配给每个任务的效果更好。您也可以通过这种方式减轻数据库的负载,通过计算每个 id,将其存储在一个列表中,然后在完成后将整个列表添加到数据库中,而不是一次只做一个 id .分块方法看起来像这样

tasks.py:

@app.task
def process_ids(items):
    for item in items:
        id = item #long complicated equation here
        database.objects(newid=id).save() # Still adding one id at a time, but you don't have to.

然后开始任务:

from tasks import process_ids

jobs = process_ids.chunks(list_of_millions_of_ids, 30) # break the list into 30 chunks. Experiment with what number works best here.
jobs.apply_async()

您可以尝试一下哪种分块大小可以获得最佳结果。您想找到一个最佳点,在此减少消息传递开销,同时保持足够小的大小,这样您最终不会让工作人员比另一个工作人员更快地完成他们的工作块,然后无所事事地等待。

【讨论】:

  • 因此,我执行“复杂的 CPU 繁重任务(可能是 3d 渲染)”的部分将自动分布式并行处理,即 1 个任务将使用所有实例可用的尽可能多的处理能力 - - 所有这些都是开箱即用的?真的吗?哇。 PS很好的答案感谢您更好地向我解释这一点。
  • @Spike 不完全是。当前编写的任务只能使用一个内核。要使单个任务使用多个核心,我们将引入threadingmultiprocessing。我们没有这样做,而是让每个 celery worker 产生与机器上可用内核一样多的任务(这在 celery 中默认发生)。这意味着在整个集群中,每个核心都可以用于处理您的list_of_million_ids,方法是让每个任务使用一个核心。因此,我们不是让单个任务使用多个内核,而是让多个任务每个都使用一个内核。这有意义吗?
  • “要让单个任务使用多个内核,我们需要引入threadingmultiprocessing”。假设我们不能将繁重的任务拆分为多个任务,您将如何使用线程或多处理来让 celery 在多个实例之间拆分任务?谢谢
  • @Tristan 这取决于任务实际在做什么。但是,在大多数情况下,我会说,如果您不能将任务本身拆分为子任务,那么您可能很难使用 multiprocessing 从任务本身内部拆分工作,因为这两种方法最终都会需要做同样的事情:将任务拆分为可以并行运行的较小任务。您实际上只是在更改进行拆分的点。
  • @PirateApp 这个问题是说你不能在 Celery 任务中使用multiprocessing inside。 Celery 本身正在使用billiardmultiprocessing fork)在单独的进程中运行您的任务。只是不允许您在其中使用multiprocessing
【解决方案2】:

在分发的世界中,您最应该记住的只有一件事:

过早的优化是万恶之源。作者:D. Knuth

我知道这听起来很明显,但在分发仔细检查之前,您使用的是最好的算法(如果存在的话......)。 话虽如此,优化分配是三件事之间的平衡:

  1. 从持久性介质写入/读取数据,
  2. 将数据从介质 A 移动到介质 B,
  3. 处理数据,

计算机的制造使您离处理单元 (3) 越近,(1) 和 (2) 就会越快、越高效。经典集群中的顺序将是:网络硬盘驱动器、本地硬盘驱动器、RAM、内部处理单元区域...... 如今,处理器变得足够复杂,可以被视为通常称为核心的独立硬件处理单元的集合,这些核心通过线程 (2) 处理数据 (3)。 想象一下,你的内核是如此之快,以至于当你用一个线程发送数据时,你使用了 50% 的计算机功率,如果内核有 2 个线程,你将使用 100%。每个内核两个线程称为超线程,您的操作系统将看到每个超线程内核有 2 个 CPU。

在处理器中管理线程通常称为多线程。 从操作系统管理 CPU 通常称为多处理。 在集群中管理并发任务通常称为并行编程。 管理集群中的依赖任务通常称为分布式编程。

那么你的瓶颈在哪里?

  • 在 (1) 中:尝试从上层(离您的处理单元较近的层,例如如果网络硬盘驱动器速度较慢首先保存在本地硬盘驱动器中)进行持久化和流式传输。
  • 在(2)中:这是最常见的一种,尽量避免分发不需要的通信数据包或压缩“即时”数据包(例如,如果 HD 很慢,则只保存“批量计算”消息并将中间结果保存在 RAM 中)。
  • 在 (3) 中:你完成了!您正在使用所有可用的处理能力。

芹菜怎么样?

Celery 是用于分布式编程的消息传递框架,它将使用代理模块进行通信 (2) 和后端模块进行持久性 (1),这意味着您将能够通过更改配置来避免大多数瓶颈(如果可能)在您的网络上并且仅在您的网络上。 首先分析您的代码以在单台计算机中实现最佳性能。 然后在集群中使用默认配置的 celery 并设置 CELERY_RESULT_PERSISTENT=True

from celery import Celery

app = Celery('tasks', 
             broker='amqp://guest@localhost//',
             backend='redis://localhost')

@app.task
def process_id(all_the_data_parameters_needed_to_process_in_this_computer):
    #code that does stuff
    return result

在执行过程中打开你喜欢的监控工具,我使用默认的rabbitMQ和celery的flower和cpus的top,你的结果将保存在你的后端。网络瓶颈的一个例子是任务队列增长如此之快以至于它们延迟执行,你可以继续更改模块或 celery 配置,如果不是你的瓶颈在其他地方。

【讨论】:

    【解决方案3】:

    为什么不使用group celery 任务呢?

    http://celery.readthedocs.org/en/latest/userguide/canvas.html#groups

    基本上,您应该将ids 划分为块(或范围)并将它们分配给group 中的一堆任务。

    对于更复杂的事情,比如聚合特定 celery 任务的结果,我已经成功地将 chord 任务用于类似目的:

    http://celery.readthedocs.org/en/latest/userguide/canvas.html#chords

    settings.CELERYD_CONCURRENCY增加到一个合理且您负担得起的数字,然后那些芹菜工人将继续以小组或和弦的形式执行您的任务,直到完成。

    注意:由于 kombu 中的一个错误,过去在大量任务中重用工作人员时会遇到问题,我不知道现在是否已修复。也许是,但如果不是,减少 CELERYD_MAX_TASKS_PER_CHILD。

    基于我运行的简化和修改代码的示例:

    @app.task
    def do_matches():
        match_data = ...
        result = chord(single_batch_processor.s(m) for m in match_data)(summarize.s())
    

    summarize 获取所有single_batch_processor 任务的结果。每个任务都在任何 Celery worker 上运行,kombu 协调它。

    现在我明白了:single_batch_processorsummarize 还必须是 celery 任务,而不是常规函数 - 否则当然不会并行化(如果不是芹菜任务)。

    【讨论】:

    • 据我了解,这会将任务拆分,但不会将芹菜并行分布式任务与多处理一起使用。即只使用所有云机器上的所有免费 CPU 能力。
    • 我不确定为什么会发生这种情况 - Celery 的工作方式就像你有一群工人,无论他们位于何处,他们甚至可能位于另一台机器上。当然,您需要有不止一名工人。 chord(将 CELERYD_CONCURRENCY 设置为数十个工作线程 == 逻辑 CPU / 硬件线程)是我在多个内核上以并行方式处理大量日志文件批次的方式。
    • 这是一个非常糟糕的代码示例。 任务do_matches 将被等待和弦阻塞。这可能会导致部分或全部死锁,因为许多/所有工作人员可能会等待子任务,而这些工作都不会完成(因为工作人员等待子任务而不是努力工作)。
    • @PrisacariDmitrii 那么正确的解决方案是什么?
    【解决方案4】:

    添加更多 celery worker 肯定会加快任务的执行速度。不过,您可能还有另一个瓶颈:数据库。确保它可以处理同时插入/更新。

    关于您的问题:您正在通过将 EC2 实例上的另一个进程分配为 celeryd 来添加芹菜工人。根据您需要的工作人员数量,您可能需要添加更多实例。

    【讨论】:

    • > 添加更多 celery worker 肯定会加快任务的执行速度。 - - 可以?所以你说的 celery 会在我所有的实例中分配一个任务,而我不必把它分开?
    • 等一下。我刚刚再次阅读了您的代码,因为它只是一项任务,这无济于事。您可以为每个 ID(或 ID 块)触发一个任务。或者您在另一个答案中遵循 John Doe 的建议。然后你可以从芹菜工人的数量中获利。是的,在这种情况下,您不需要做太多事情。只需确保工作人员使用相同的队列即可。
    猜你喜欢
    • 2018-11-13
    • 2018-06-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-01-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多