【问题标题】:Adding new column in pyspark dataframe在 pyspark 数据框中添加新列
【发布时间】:2021-07-09 12:45:20
【问题描述】:

我正在尝试将新记录时区添加到我的 pysaprk 数据帧

from timezonefinder import TimezoneFinder
tf = TimezoneFinder()
df = df.withColumn("longitude",col("longitude").cast("float"))
df = df.withColumn("Latitude",col("Latitude").cast("float"))
df = df.withColumn("timezone",tf.timezone_at(lng=col("longitude"), lat=col("Latitude")))

我遇到了错误。

ValueError: Cannot convert column into bool: please use '&' for 'and', '|' for 'or', '~' for 'not' when building DataFrame boolean expressions.

Timezonefinder 库用于通过传递地理坐标来查找时区。

Latitude, longitude = 20.5061, 50.358
tf.timezone_at(lng=longitude, lat=Latitude)
 -- 'Asia/Riyadh'

【问题讨论】:

    标签: python apache-spark pyspark user-defined-functions


    【解决方案1】:

    您需要使用 UDF 将列传递给 Python 函数:

    import pyspark.sql.functions as F
    
    @F.udf('string')
    def tfUDF(lng, lat):
        from timezonefinder import TimezoneFinder
        tf = TimezoneFinder()
        return tf.timezone_at(lng=lng, lat=lat)
    
    df = df.withColumn("longitude", F.col("longitude").cast("float"))
    df = df.withColumn("Latitude", F.col("Latitude").cast("float"))
    df = df.withColumn("timezone", tfUDF(F.col("longitude"), F.col("Latitude")))
    
    df.show()
    +--------+---------+-----------+
    |Latitude|longitude|   timezone|
    +--------+---------+-----------+
    | 20.5061|   50.358|Asia/Riyadh|
    +--------+---------+-----------+
    

    【讨论】:

    • 感谢它的工作。只是想知道是否有任何替代使用UDF。还有使用 UDF 的性能如何。我是使用 Pyspark 的绝对初学者。
    • 您可能需要为 Spark SQL 中本机不可用的任何函数使用 UDF。它的性能通常比原生 Spark 函数差,但您可以考虑使用 pandas UDF 和 pyarrow 来提高 UDF 性能。
    猜你喜欢
    • 1970-01-01
    • 2015-11-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多