【问题标题】:Dask HighLevelGraph short circuit computingDask HighLevelGraph 短路计算
【发布时间】:2019-10-15 13:12:45
【问题描述】:

我正在尝试获取一个 DataFrame ddf 并返回一个与 ddf 相同的新 DataFrame,除非 ddf 有一个空分区,它应该指向最近的非空组件。例如,如果ddf 具有分区[P1, P2, P3, P4, P5, P6],其中P2P3P6 是空的Pandas DataFrame,那么它返回以下Dask DataFrame:[P1, P1, P1, P4, P5, P5]。我的代码是

name = 'prev-nonempty-' + tokenize(ddf)
meta = ddf._meta
dsk = dict()
def helper(A, B):
  return B if A.empty else A
dsk[(name, 0)] = (helper, (ddf._name, 0), None)
for i in range(1, len(ddf.divisions)-1):
    dsk[(name, i)] = (helper, (ddf._name, i), (name, i-1))
graph = HighLevelGraph.from_collections(name, dsk, dependencies=[ddf])
return new_dd_object(graph, name, meta, ddf.divisions)

我的问题是是否有办法在 Dask HighLevelGraphs 中进行短路计算,以便如果找到非空分区,第 i 个分区的计算会提前停止。

上面写着here

在像(add, 'x', 'y') 这样的情况下,像add 这样的函数接收具体的值而不是键。 Dask 调度程序将键(如 xy)替换为其计算值(如 12在调用 add 函数之前

这表明您不能将其短路,但也许我可以使用更复杂的 Dask 调度程序技巧?

【问题讨论】:

    标签: dataframe dask short-circuiting


    【解决方案1】:

    不,标准任务图无法做到这一点。但是,您可以将此逻辑融入您的函数本身。

    def func(accumulator, new_data):
        if is_done(accumulator):
            return accumulator 
    

    所以你仍然会跑完所有的任务,但是在你满足你的条件之后它们会非常快。

    您也可以考虑使用 Dask Futures,但这个级别有点低。 https://docs.dask.org/en/latest/futures.html

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-06-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-12-10
      • 1970-01-01
      • 2021-08-26
      • 1970-01-01
      相关资源
      最近更新 更多