【问题标题】:PySpark: How to apply a Python UDF to PySpark DataFrame columns?PySpark:如何将 Python UDF 应用于 PySpark DataFrame 列?
【发布时间】:2020-05-21 23:32:06
【问题描述】:

我有一个 PySpark DataFrame,其中包含两组纬度、经度坐标。我正在尝试计算给定行的每组坐标之间的 Haversine 距离。我正在使用我在网上找到的以下haversine()。问题是它不能应用于列,或者至少我不知道这样做的语法。有人可以分享语法或指出更好的解决方案吗?

from math import radians, cos, sin, asin, sqrt

def haversine(lat1, lon1, lat2, lon2):
    """
    Calculate the great circle distance between two points 
    on the earth (specified in decimal degrees)
    """
    # convert decimal degrees to radians 
    lon1, lat1, lon2, lat2 = map(radians, [lon1, lat1, lon2, lat2])
    # haversine formula 
    dlon = lon2 - lon1 
    dlat = lat2 - lat1 
    a = sin(dlat/2)**2 + cos(lat1) * cos(lat2) * sin(dlon/2)**2
    c = 2 * asin(sqrt(a)) 
    # Radius of earth in miles is 3,963; 5280 ft in 1 mile
    ft = 3963 * 5280 * c
    return ft

我知道上面的 haversine() 函数有效,因为我使用数据框中的一些纬度/经度坐标对其进行了测试并得到了合理的结果:

haversine(-85.8059, 38.250134, 
          -85.805122, 38.250098)
284.1302325439314

当我在 PySpark 数据框中将示例坐标替换为对应于 lat/lons 的列名时,出现错误。我尝试了以下代码,试图创建一个新列,其中包含计算的Haversine 距离(以英尺为单位):

df.select('id', 'p1_longitude', 'p1_latitude', 'p2_lon', 'p2_lat').withColumn('haversine_dist', 
                           haversine(df['p1_latitude'],
                                    df['p1_longitude'],
                                    df['p2_lat'],
                                    df['p2_lon']))
.show()

但我得到了错误:

必须是实数,而不是 Column Traceback(最近一次调用最后一次):
文件“”,第 8 行,haversine TypeError: must be real number, 不是列

这向我表明我必须以某种方式迭代地将我的 hasrsine 函数应用于我的 PySpark DataFrame 的每一行,但我不确定这个猜测是否正确,即使是这样,我也不知道该怎么做。顺便说一句,我的 lat/lons 是浮点类型。

【问题讨论】:

  • 尽可能避免使用 udf-s - spark 无法优化它们,它们的性能比使用内置函数或 SQL 差几倍。 SQL 中提供了所有三角函数(我不知道 radians)...因此,如果您的性能坦克尝试用 SQL 表达式重写它,看看会发生什么

标签: python apache-spark pyspark apache-spark-sql


【解决方案1】:

当您可以使用 Spark 内置函数时,不要使用 UDF,因为它们的性能通常较低。

这是一个仅使用与您的函数相同的 Spark SQL 函数的解决方案:

from pyspark.sql.functions import col, radians, asin, sin, sqrt, cos

df.withColumn("dlon", radians(col("p2_lon")) - radians(col("p1_longitude"))) \
  .withColumn("dlat", radians(col("p2_lat")) - radians(col("p1_latitude"))) \
  .withColumn("haversine_dist", asin(sqrt(
                                         sin(col("dlat") / 2) ** 2 + cos(radians(col("p1_latitude")))
                                         * cos(radians(col("p2_lat"))) * sin(col("dlon") / 2) ** 2
                                         )
                                    ) * 2 * 3963 * 5280) \
  .drop("dlon", "dlat")\
  .show(truncate=False)

给予:

+-----------+------------+----------+---------+------------------+
|p1_latitude|p1_longitude|p2_lat    |p2_lon   |haversine_dist    |
+-----------+------------+----------+---------+------------------+
|-85.8059   |38.250134   |-85.805122|38.250098|284.13023254857814|
+-----------+------------+----------+---------+------------------+

您可以找到可用的 Spark 内置函数 here。

【讨论】:

  • 这看起来很有希望。但是,我收到错误消息:<console>:33: error: value ** is not a member of org.apache.spark.sql.Column
  • 嗯,这很奇怪!上面的代码在 Spark 2.4 中可以正常运行...不知道为什么会出现该错误,但或者,您可以使用 Spark SQL 函数 power 而不是 Python 运算符 ** (a ** b 相当于 power(a, b) )。 .
猜你喜欢
  • 1970-01-01
  • 2016-11-24
  • 1970-01-01
  • 2019-10-02
  • 2017-03-16
  • 1970-01-01
  • 2021-03-04
  • 2016-08-24
  • 2020-12-12
相关资源
最近更新 更多