【问题标题】:Multiprocessing and Dask多处理和 Dask
【发布时间】:2021-03-17 12:27:14
【问题描述】:

场景:我想对一个文件(~100 MB)执行简单的计算。我有数千个这样的文件。对于此示例,将计算视为“计算行数”。

Dask:我正在使用 Dask 来并行化读取和计算。

工作原理:如果每个文件都非常大(45 GB),那么在我的集群(12 台使用 PBS 的计算机)中专用 1 台计算机(70GB 内存)来读取和计算是有意义的。

from dask_jobqueue import PBSCluster
cluster = PBSCluster(cores=1,processes=1,memory='70GB',queue='extra',project='abcd',walltime='12:00:00',interface='ib0',local_directory='/home/abcd/dask_workers/')
cluster.scale(jobs=12)
print(cluster.job_script())
from dask.distributed import Client, progress
client = Client(cluster)

问题:当文件很小(~100 MB)时,我希望集群中的 1 台计算机并行读取和计算多个文件。

问题:哪个参数将允许集群中的 1 台计算机并行进行多次读取和计算。

【问题讨论】:

    标签: python dask dask-distributed


    【解决方案1】:

    指定一些任意资源怎么样?例如。在创建集群时指定resources={'compute_load': 12},给每个worker 12个单位的资源compute_load

    处理重负载时,分配.compute(resources={'compute_load': 12}),这样每个worker一次只承担一个任务。当负载较轻时,使用.compute(resources={'compute_load': 1}),每个worker一次会占用12个任务。另见this answer

    【讨论】:

    • 我认为这行不通。我是这样设置的:cluster = PBSCluster(resources={'compute_load':12},cores=1,memory='10 GB',queue='extra',project='abcd',walltime='10:00:00',interface='ib0',local_directory='/home/abcd/dask_workers/')cluster.scale(12)client = Client(cluster)futures = [client.submit(functionToProcess,inputFile,resources={'compute_load':1}) for eachFile in intake]
    • 这看起来不错,但我认为您想使用 futures = [client.submit(functionToProcess,inputFile,resources={'compute_load':12}) for eachFile in intake] 处理大量文件。
    • 还请注意,每个工作人员只有一个核心,他们将无法并行执行任何操作...
    • 此设置适用于轻负载。每个工作人员有两个处理器,每个处理器有 6 个内核。
    • hmm,cores=1 表示应该为每个worker分配一个核心......
    猜你喜欢
    • 2020-03-02
    • 2020-07-13
    • 1970-01-01
    • 2021-06-08
    • 1970-01-01
    • 1970-01-01
    • 2017-08-09
    • 1970-01-01
    • 2022-01-19
    相关资源
    最近更新 更多