【问题标题】:Matching groups of data with Dask使用 Dask 匹配数据组
【发布时间】:2020-05-03 02:32:20
【问题描述】:

问题和问题陈述

我有来自两个来源的数据。每个源都包含由ID 列、坐标和属性标识的组。我想通过首先匹配这些组来处理这些数据,然后在这些组中找到最近的邻居,然后研究来自不同来源的属性如何在邻居之间进行比较。我对自己的学习挑战是如何使用并行处理来处理这些数据。

问题是:“使用 Dask 进行并行处理,处理此类数据的最简单、最直接的方法可能是什么?”

到目前为止的背景和我的解决方案

数据在 CSV 文件中,如下面的虚拟数据(真实文件在 100 MiB 范围内):

source1.csv:
ID,X_COORDINATE,Y_COORDINATE,ATTRIB1,PARAM1
B,-63802.84728184705,-21755.63629150563,3,36.136464492674556
B,-63254.41147034371,405.6973789009853,1,18.773534321367528
A,-9536.906537069272,32454.934987740824,0,14.043507555168809
A,15250.802157581298,-40868.390394552596,0,6.680542212635015
source2.csv:
ID,X_COORDINATE,Y_COORDINATE,ATTRIB1,PARAM1
B,-6605.150024790153,39733.35763934722,3,5.599467583303852
B,53264.28797042654,24647.24183964514,0,27.938127686688162
A,6690.836682554512,34643.0606728128,0,10.02914141165683
A,15243.16,-40954.928,0,18.130371948545935

我想做的是

  1. 将数据加载到数据帧中
  2. 按 ID 列将它们分组
  3. 对于source1source2 中的每个组,让我们调用每个组中的子数据帧source1_subsource2_sub
    • 根据列 X_COORDINATE 和 Y_COORDINATE 构造 kdtree 对象 k1k2
  4. 对于每对对象(k1, k2)
    • 为树木找到最近的邻居
    • 构造三个数据帧:
      • matches_sub:包含source1_subsource2_sub 中的匹配行
      • source1_sub_onlysource1_sub 中不匹配的行
      • source2_sub_onlysource2_sub 中不匹配的行
  5. 将所有matches_subsource1_sub_onlysource2_sub_only数据帧连接成三个数据帧:matchessource1_onlysource2_only
  6. 分析这些数据帧

这是一个应该很好地并行化的问题,因为每对组都独立于其他组对。我决定使用scipy.spatial.cKDTree 进行实际的坐标匹配,但困难在于它对原始 numpy 数组的索引进行操作,这与访问 Dask 数组的方式并不直接兼容。至少这是我的理解。

我的第一次徒劳的尝试真的很尴尬

  1. 尝试使用两个 Dask 数据帧,对齐它们并找到匹配项。这非常缓慢且难以理解。
  2. 使用 Dask Dataframe 读取数据并使用 Dask Bag 进行处理。这稍微不那么复杂,但仍不能令人满意。

【问题讨论】:

  • 我编辑了问题以重新表述实际想要做的事情。

标签: python pandas dataframe dask


【解决方案1】:

回答我自己,我能想到的最简单的方法是

  1. 使用 dask.dataframe.read_csv 将来自源 1 和 2 的数据读入数据帧 df_source1df_source2
  2. 在读取时,将新列 SOURCE 分配给这些数据帧,以识别来源。现在我感兴趣的组由IDSOURCE 列指定。这可用于分组。
  3. 将这些数据帧连接到新的数据帧df = dd.concat([df_source1, df_source2], axis=0)
  4. IDSOURCE 列对数据进行分组,并使用apply 查找匹配项。
  5. 分析数据。
  6. 完成。

类似的东西:

import dask.dataframe as dd

import pandas as pd
import numpy as np

from scipy.spatial import cKDTree

def find_matches(x):
    x_by_source = x.groupby(['SOURCE'])

    grp1 = x_by_source.get_group(1)
    grp2 = x_by_source.get_group(2)

    tree1 = cKDTree(grp1[['X_COORDINATE', 'Y_COORDINATE']])
    tree2 = cKDTree(grp2[['X_COORDINATE', 'Y_COORDINATE']])
    neighbours = tree1.query_ball_tree(tree2, r=70000)
    matches = np.array([[n,k] for (n, j) in enumerate(neighbours) if j != [] for k in j])

    indices1 = grp1.index[matches[:,0]]
    indices2 = grp2.index[matches[:,1]]

    m1 = grp1.loc[indices1]
    m2 = grp2.loc[indices2]

    # arrange matches side by side
    res = pd.concat([m1, m2], ignore_index=True, axis=1)

    return(res)

df_source1 = dd.read_csv('source1.csv').assign(SOURCE = 1)
df_source2 = dd.read_csv('source2.csv').assign(SOURCE = 2)

df = dd.concat([df_source1, df_source2], axis=0)

meta = pd.DataFrame(columns=np.arange(0, 2*len(df.columns)))

result = (df.groupby('ID')
    .apply(find_matches, meta=meta)
    .persist()
)

# Proceed with further analysis

【讨论】:

  • 经过两周的沉思后,我得出了这个答案。我意识到我原来的方法太复杂了。我在这里回答,以防它也有利于其他人。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-02-12
  • 1970-01-01
  • 2017-08-20
  • 2021-09-13
相关资源
最近更新 更多