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 实现映射减少范例,例如将一个巨大的数组分成较小的块,进行一些处理并将结果返回
我是否缺少一些关键配置?
Map-reduce 范式不应该更快,而是应该更好地扩展。
与实现相同计算的本地运行作业相比,MR 作业始终存在开销:进程调度、通信、洗牌等。
您的基准测试并不相关,因为 MR 和本地运行都是方法,具体取决于数据集大小。在某些时候,您会从本地运行方法切换到 MR 方法,因为您的数据集对于一个节点来说太大了。