【发布时间】:2019-09-20 00:04:05
【问题描述】:
我有一个大数据集(5000 万行),我需要在其中进行一些逐行计算,例如获取两个集合的交集(每个集合在不同的列中)
例如
col_1:{1587004, 1587005, 1587006, 1587007}
col_2:{1587004, 1587005}
col_1.intersection(col_2) = {1587004, 1587005}
这适用于我的虚拟数据集 (100 000) 行。 但是,当我尝试与实际相同时,内存耗尽
我的编码使用 pandas 1:1 将其移植到 dask 不起作用 NotImplementedError: 系列 getitem 仅支持其他具有匹配分区结构的系列对象
到目前为止,使用 map_partitions 并没有奏效
工作代码:
df["intersection"] = [col_1.intersection(col_2) for col_1,col2 in zip(df.col_1,df.col_2)]
用 dask df 替换 pandas df 在未实现的错误中运行:
ddf["intersection"] = [col_1.intersection(col_2) for col_1,col2 in zip(df.col_1,df.col_2)]
使用 map_partions “有效”,但我不知道如何将结果分配给现有的 ddf
def intersect_sets(df, col_1, col_2):
result = df[col_1].intersection(df[col_2])
return result
newCol = ddf.map_partitions(lambda df : df.apply(lambda series: intersect_sets(series,"col_1","col_2"),axis=1),meta=str).compute()
只是在做:
ddf['result'] = newCol
导致:
ValueError: 并非所有分区都是已知的,无法对齐分区。请使用set_index设置索引。
更新: 重置索引会消除错误,但是包含交叉点的列不再与其他两列匹配。顺序好像乱了……
ddf2 = ddf.reset_index().set_index('index')
ddf2 ['result'] = result
我希望有一个包含以下列的 dask 数据框
col_1:{1587004, 1587005, 1587006, 1587007}
col_2:{1587004, 1587005}
col_3:{1587004, 1587005}
不仅感谢完美的解决方案,而且对 map_partitions 如何工作的一些见解已经对我有很大帮助:)
更新: 感谢 M.Rocklin,我想通了。 对于将来我或其他人在这个问题上磕磕绊绊:
ddf = ddf.assign(
new_col = ddf.map_partitions(
lambda df : df.apply(
lambda series:intersect_sets(
series,"col_1","col_2"),axis=1),meta=str)
)
df = ddf.compute()
【问题讨论】: