【发布时间】:2019-10-15 13:12:45
【问题描述】:
我正在尝试获取一个 DataFrame ddf 并返回一个与 ddf 相同的新 DataFrame,除非 ddf 有一个空分区,它应该指向最近的非空组件。例如,如果ddf 具有分区[P1, P2, P3, P4, P5, P6],其中P2、P3 和P6 是空的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 调度程序将键(如x和y)替换为其计算值(如1和2)在调用add函数之前。
这表明您不能将其短路,但也许我可以使用更复杂的 Dask 调度程序技巧?
【问题讨论】:
标签: dataframe dask short-circuiting