【发布时间】:2020-12-06 16:58:39
【问题描述】:
如何将以下函数移植到 dask 以使其并行化?
from time import sleep
from dask.distributed import Client
from dask import delayed
client = Client(n_workers=4)
from tqdm import tqdm
tqdm.pandas()
# linear
things = [1,2,3]
_x = []
_y = []
def my_slow_function(foo):
sleep(2)
x = foo
y = 2 * foo
assert y < 5
return x, y
for foo in tqdm(things):
try:
x_v, y_v = my_slow_function(foo)
_x.append(x_v)
if y_v is not None: _y.append(y_v)
except AssertionError:
print(f'failed: {foo}')
X = _x
y = _y
print(X)
print(y)
我特别不确定如何处理延迟期货中的状态和故障。
目前我只有:
from dask.diagnostics import ProgressBar
ProgressBar().register()
@delayed(nout=2)
def my_slow_function(foo):
sleep(2)
x = foo
y = 2 * foo
assert y < 5
return x, y
for foo in tqdm(things):
try:
x_v, y_v = delayed(my_slow_function(foo))
_x.append(x_v)
if y_v is not None: _y.append(y_v)
except AssertionError:
print(f'failed: {foo}')
X = _x
y = _y
print(X)
print(y)
delayed(sum)(X).compute()
但是:
- try/except 不再有效。 IE。不再捕获异常
- 我有 2 个延迟结果列表,但没有 2 个计算值列表
- 对于这 2 个列表,我不确定如何在不计算两次结果的情况下执行
compute
- 对于这 2 个列表,我不确定如何在不计算两次结果的情况下执行
编辑
futures = client.map(my_slow_function, things)
results = client.gather(futures)
由于不再处理异常而明显失败 - 但到目前为止,我不确定从 dask 中捕获它们的正确方法是什么。
How to prevent dask client from dying on worker exception? 可能类似
【问题讨论】:
标签: python dask future dask-delayed