【问题标题】:How do the batching instructions of Dask delayed best practices work?Dask 延迟最佳实践的批处理指令如何工作?
【发布时间】:2019-10-20 12:48:53
【问题描述】:

我想我遗漏了一些东西(仍然是 Dask Noob),但我正在尝试批处理建议以避免从这里执行过多的 Dask 任务:

https://docs.dask.org/en/latest/delayed-best-practices.html

并且不能让它们工作。 这是我尝试过的:

import dask

def f(x):
    return x*x

def batch(seq):
    sub_results = []
    for x in seq:
        sub_results.append(f(x))
    return sub_results

batches = []
for i in range(0, 1000000000, 1000000):
 result_batch = dask.delayed(batch, range(i, i + 1000000))
 batches.append(result_batch)

批次现在包含延迟对象:

batches[:3]

[Delayed(range(0, 1000000)),
 Delayed(range(1000000, 2000000)),
 Delayed(range(2000000, 3000000))]

但是当我计算它们时,我会得到批处理函数指针(我认为??):

results = dask.compute(*batches)
results[:3]

(<function __main__.batch(seq)>,
 <function __main__.batch(seq)>,
 <function __main__.batch(seq)>)

我有两个问题:

  1. 这真的应该如何运行,因为它似乎与 Best practices 页面的第一行相矛盾,它说 not 像 delayed(f(x)) 一样运行它,因为那样会立即运行而不是偷懒。

  2. 如何获得上述批量运行的结果?

【问题讨论】:

  • 哇,没有任何评论或答案的热门问题徽章,这是第一次。 ;)
  • 当然uncommon.
  • ..为我.......

标签: python dask dask-delayed


【解决方案1】:

您的代码似乎缺少一对括号。不确定这是否是一个错字 (???)。

根据文档中的示例,我认为您想要

result_batch = dask.delayed(batch)(range(i, i + 1000000))

我将batch, ran... 替换为batch)(ran...,因为应该延迟调用batch() 函数。

答案

  1. 修正错字后,您的代码对我来说可以正常工作 - 现在计算将被延迟。关于文档开头所写的内容 - 使用 dask.delayed 包装的内容很重要。使用dask.delayed( batch(range(i, i + 1000000)) ) 对函数batch(...) 的调用不会被延迟,因此它会立即运行。这是因为函数的输出已经包裹在dask.delayed 中,所以输出(结果)会被延迟,这不是我们想要的工作流程。然而,dask.delayed(batch)(range(i, i + 1000000)) 延迟了对函数的调用(因为这里,dask.delayed 包装了函数本身)。我相信这就是文档在最佳实践部分开始时想要表达的意思。
  2. 同样,在修正错字后,您的代码可以按预期运行,并将 longy 输出打印到屏幕上。

【讨论】:

  • 是的,有道理。我只是从他们的文档中复制了这个,当时没有理解,因为那里也是错误的:github.com/dask/dask/commit/…
  • 啊,我没想到要像问题那样深入。这就解释了。
猜你喜欢
  • 1970-01-01
  • 2017-04-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多