如何解决Dask DataFrame的“not all divisions are known”报错?
Dask分组过滤后调用
head()报错的解决方案 问题场景
你写了这段代码来过滤Dask DataFrame里出现次数超过15次的ngram:
df = df.set_index('ngram') sizes = df.groupby('ngram').size() df = df[sizes > 15]
结果调用df.head(15)的时候直接触发错误:
ValueError: Not all divisions are known, can't align partitions. Please use
set_indexto set the index.
你检查分区状态后发现:
>>> df.known_divisions # 返回True >>> sizes.known_divisions # 返回False
目前你用了临时办法:把sizes落地(compute()或者存成CSV再读回来),但这个方案有点绕,其实有更高效的解决方式。
问题根因
核心矛盾就是分区信息不匹配:
- 你给
df设置索引后,它的分区是明确的(known_divisions=True); - 但
groupby.size()生成的sizes对象,Dask无法自动推断其分区规则,所以known_divisions=False; - 当用这个分区未知的
sizes去过滤df时,Dask没办法对齐两个对象的分区逻辑,最终导致报错。
更高效的修复方案
不用落地数据,直接给sizes同步df的分区信息即可:
df = df.set_index('ngram') sizes = df.groupby('ngram').size() # 关键一步:复用df的分区规则来设置sizes的分区 sizes = sizes.set_index(sizes.index, divisions=df.divisions) # 现在再执行过滤就不会有问题 df = df[sizes > 15] # 调用head()完全正常 df.head(15)
执行完这步后,sizes.known_divisions会变成True,分区对齐的问题直接解决。
为什么你的临时办法能生效?
- 用
sizes.compute()会把它转成Pandas Series,Pandas没有分区概念,Dask会把这个本地对象广播到每个分区完成过滤; - 存成CSV再读回来时,新生成的Dask对象会默认生成已知分区,所以也能解决问题,但这两种方法都会额外消耗内存或IO资源,效率远不如直接同步分区高。
内容的提问来源于stack exchange,提问作者Dominykas Mostauskis
相关产品推荐
相关产品推荐

