【问题标题】:Dask - map_partitionDask - map_partition
【发布时间】:2021-12-22 13:30:16
【问题描述】:

我有一个带有纬度和经度集(约 32m 条记录)的 Dask DataFrame。我正在尝试使用如下函数计算纬度/经度之间的距离:

import numpy as np
from geopy import distance

def calc_distance(df, lat_col_name_1, lon_col_name_1, lat_col_name_2, lon_col_name_2):
if df[lat_col_name_1] != np.nan and df[lon_col_name_1] != np.nan and df[lat_col_name_2] != np.nan and df[lon_col_name_2] != np.nan:
    return distance.distance((df[lat_col_name_1], df[lon_col_name_1]), (df[lat_col_name_2], df[lon_col_name_2])).miles
else:
    return np.nan 

我尝试使用 map_partitions 调用此函数(以创建索引和距离的 DataFrame 以及使用 assign 调用 map_paritions。我想使用 assign 以避免将 DataFrame 重新连接在一起(似乎很昂贵)。它确实不像 np.nan 检查。我得到了一个

ValueError:Series 的真值不明确。使用 a.empty, a.bool()、a.item()、a.any() 或 a.all()。

我有零纬度/经度的记录,所以我需要在计算距离时考虑到这一点。

使用 map_partitions

distance = big_df.map_partitions(calc_distance, 
                                    lat_col_name_1='latitude_1', 
                                    lon_col_name_1='longitude_1', 
                                    lat_col_name_2='latitude_2', 
                                    lon_col_name_2='longitude_2', 
                                    meta={'distance': np.float64})

使用 map_partitions 和分配

def calc_distance_miles(lat1, lon1, lat2, lon2):
    if lat1 != np.nan and lon1 != np.nan and lat2 != np.nan and lon2 != np.nan:
        return distance.distance((lat1, lon1), (lat2, lon2)).miles
    else:
        return np.nan
    

big_df = big_df.map_partitions(lambda df: df.assign(
    distance=calc_distance_miles(df['latitude_1'], df['longitude_1'], df['latitude_2'], df['longitude_2'])
), meta={'distance': np.float64}
)

【问题讨论】:

  • 小心使用np.nan 的布尔运算符。 NaN 从不等于任何东西。请注意,np.nan != np.nan 的计算结果为 True。所以你的测试没有做任何事情。相反,在 DataFrame 和 Series 上使用 pd.isnull()isnull 方法。见the pandas docs on working with missing data

标签: python dataframe dask dask-distributed dask-dataframe


【解决方案1】:

map_partitions 不像 df.apply - 函数 calc_distance 是通过 dask.dataframe 的分区调用的,该分区的类型为 pd.DataFrame。

因此,df[lat_col_name_1] 是一个系列,df[lat_col_name_1] != np.nan 是一个布尔系列(它将始终返回此错误 - 参见例如 Truth value of a Series is ambiguous. Use a.empty, a.bool(), a.item(), a.any() or a.all())。

有比按元素计算距离更快的数组方法,但与您尝试做的类似的 dask.dataframe 是使用 map_partitions 然后应用:

def calc_distance(series, lat_col_name_1, lon_col_name_1, lat_col_name_2, lon_col_name_2):
    if series[
        [lat_col_name_1, lon_col_name_1, lat_col_name_2, lon_col_name_2]
    ].notnull().all():

        return distance.distance(
            (series[lat_col_name_1], series[lon_col_name_1]),
            (series[lat_col_name_2], series[lon_col_name_2]),
        ).miles

    else:
        return np.nan 

def calc_distance_df(df, **kwargs):
    return df.apply(calc_distance, axis=1, **kwargs)

distances = big_df.map_partitions(
    calc_distance_df,
    meta=np.float64,
    lat_col_name_1=lat_col_name_1,
    lon_col_name_1=lon_col_name_1,
    lat_col_name_2=lat_col_name_2,
    lon_col_name_2=lon_col_name_2,
)

【讨论】:

  • 感谢您的解释。不幸的是,这并不完全奏效。移植时出现此错误:“[Index(['latitude_1', 'longitude_1', 'latitude_2',\n 'longitude_2'],\n dtype='object')] 中没有 [index] " 你能解释一下你所说的比计算元素更快的方法是什么意思吗?
  • oops - 抱歉,我使用了错误的轴参数。应该是axis = 1,而不是axis = 0。再试一次。
  • 使用数组在 pandas 中执行任何类型的操作都明显更快。您正在遍历数据框中的各个行(pd.apply 只是一个 for 循环),然后计算每行中元素之间的距离。您可以使用矢量化算法一次计算数据帧中所有行的每行中的点之间的距离 - 这将明显更快。不确定 geopy 是否支持这一点,但您可以使用 geopandas 或其他矢量化库。
  • 谢谢迈克尔。我现在正在测试。我最初在 pandas 中循环使用块大小,但运行时间太长(32m 记录需要几个小时)。我使用了以下内容:chunk['distance_miles'] = np.vectorize(calc_distance_miles)(chunk['point1'], chunk['point2']) 纬度/经度元组的点在哪里。这比应用函数快得多,但仍然太慢。这就是矢量化的意思吗?
  • 实际上 np.vectorize 是相似的,因为它只是循环遍历元素。我的意思是同时处理整个向量,如这个问题:vectorizing haversine distance calculation in python
猜你喜欢
  • 2021-03-13
  • 1970-01-01
  • 2017-06-18
  • 2018-05-04
  • 2023-04-04
  • 1970-01-01
  • 2019-01-07
  • 1970-01-01
  • 2017-06-02
相关资源
最近更新 更多