你的目标是:
- 将您的工作分配给多台机器(分布式
计算/分布式并行处理)
- 将给定机器上的工作分配给所有 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()
您可以尝试一下哪种分块大小可以获得最佳结果。您想找到一个最佳点,在此减少消息传递开销,同时保持足够小的大小,这样您最终不会让工作人员比另一个工作人员更快地完成他们的工作块,然后无所事事地等待。