使用 celery 在简单任务中表现不佳

Vij*_*ddi 5 python rabbitmq celery

我目前在执行以下用例时面临性能不佳的问题:

我有两个文件 -task.py

# tasks.py
from celery import Celery

app = Celery('tasks', broker='pyamqp://guest@localhost//', backend='rpc://',worker_prefetch_multiplier=1)

@app.task
def task(array_of_elements):
    return [x ** 2 for x in array_of_elements]
Run Code Online (Sandbox Code Playgroud)

和运行.py

# run.py
from celery import group
from itertools import chain, repeat
from tasks import task
import time

def grouper(n, iterable, padvalue=None):
    return zip(*[chain(iterable, repeat(padvalue, n-1))]*n)

def fun1(x):
    return x ** 2

if __name__ == '__main__':
    start = time.time()
    items = [list(x) for x in grouper(10000, range(10000))]
    x = group([task.s(item) for item in items])
    r = x.apply_async()
    d = r.get()
    end = time.time()
    print(f'>celery: {end-start} seconds')

    start = time.time()
    res = [fun1(x) for x in range(10000)]
    end = time.time()
    print(f'>normal: {end-start} seconds')
Run Code Online (Sandbox Code Playgroud)

当我尝试运行 celery 时: celery -Atasksworker --loglevel=info

并尝试运行:

python run.py
Run Code Online (Sandbox Code Playgroud)

这是我得到的输出:

>celery: 0.19174742698669434 seconds
>normal: 0.004475116729736328 seconds
Run Code Online (Sandbox Code Playgroud)

我不知道为什么芹菜的性能较差?

我试图了解如何使用 celery 实现映射减少范例,例如将一个巨大的数组分成较小的块,进行一些处理并将结果返回

我是否缺少一些关键配置?

Set*_*top 2

Map-reduce 范式不应该更快,而是应该更好地扩展。

与实现相同计算的本地运行作业相比,MR 作业始终存在开销:进程调度、通信、洗牌等。

您的基准测试并不相关,因为 MR 和本地运行都是方法,具体取决于数据集大小。在某些时候,您会从本地运行方法切换到 MR 方法,因为您的数据集对于一个节点来说太大了。