【发布时间】:2017-08-18 01:54:23
【问题描述】:
dask.compute(...) 应该是一个阻塞调用。但是,当我嵌套了 dask.compute,而内部的执行 I/O(如 dask.dataframe.read_parquet)时,内部的 dask.compute 不会阻塞。这是一个伪代码示例:
import dask, distributed
def outer_func(name):
files = find_files_for_name(name)
df = inner_func(files).compute()
# do work with df
return result
def inner_func(files):
tasks = [ dask.dataframe.read_parquet(f) for f in files ]
tasks = dask.dataframe.concat(tasks)
return tasks
client = distributed.Client(scheduler_file=...)
results = dask.compute([ dask.delay(outer_func)(name) for name in names ])
如果我启动 2 个工人,每个工人有 8 个进程,例如:
dask-worker --scheduler-file $sched_file --nprocs 8 --nthreads 1
,那么我预计最多 2 x 8 个并发的 inner_func 正在运行,因为 inner_func(files).compute() 应该是阻塞的。然而,我观察到的是,在一个工作进程中,一旦它开始 read_parquet 步骤,可能会有另一个 inner_func(files).compute() 开始运行。所以最后可能有多个inner_func(files).compute()在运行,有时会导致内存不足的错误。
这是预期的行为吗?如果是这样,是否有任何方法可以强制每个工作进程执行一个 inner_func(files).compute()?
【问题讨论】:
-
这里似乎有点混乱。 dask.dataframe 创建惰性对象,在同样延迟/计算的函数中创建/计算这些对象是不正常的。考虑一下,这个函数被发送给一个工作者:你希望计算发生在哪里?
-
这个例子中的嵌套是真实世界数据流恕我直言的典型。使用像 dask DataFrame 这样的分布式数据结构实际上并不总是可行/可取的,这样我们就可以避免这种嵌套。因为 dask DataFrame API 比 pandas 小,并且因为保持一个有效的串行代码版本非常重要。从我所看到的,inner_func 似乎在 dask-worker 进程内的多个线程中运行,但我只使用例如:dask-worker --scheduler-file sched.json --nprocs 3 --nthreads 1 为每个工作者指定一个线程--local-directory /tmp/
标签: python dask dask-distributed dask-delayed