【问题标题】:spark sql distance to nearest holidayspark sql到最近假期的距离
【发布时间】:2016-11-22 21:56:48
【问题描述】:

在熊猫中我有一个类似于

的功能
indices = df.dateColumn.apply(holidays.index.searchsorted)
df['nextHolidays'] = holidays.index[indices]
df['previousHolidays'] = holidays.index[indices - 1]

计算到最近假​​期的距离并将其存储为新列。

searchsorted http://pandas.pydata.org/pandas-docs/version/0.18.1/generated/pandas.Series.searchsorted.html 是 pandas 的一个很好的解决方案,因为这给了我下一个假期的索引,而没有很高的算法复杂度 Parallelize pandas apply 例如这种方法比并行循环要快得多。

如何在 spark 或 hive 中实现这一点?

【问题讨论】:

    标签: sorting apache-spark hive apache-spark-sql datediff


    【解决方案1】:

    这可以使用聚合来完成,但这种方法比 pandas 方法具有更高的复杂性。但是您可以使用 UDF 实现类似的性能。它不会像 pandas 那样优雅,但是:

    假设这个假期数据集:

    holidays = ['2016-01-03', '2016-09-09', '2016-12-12', '2016-03-03']
    index = spark.sparkContext.broadcast(sorted(holidays))
    

    以及数据框中 2016 年日期的数据集:

    from datetime import datetime, timedelta
    dates_array = [(datetime(2016, 1, 1) + timedelta(i)).strftime('%Y-%m-%d') for i in range(366)]
    from pyspark.sql import Row
    df = spark.createDataFrame([Row(date=d) for d in dates_array])
    

    UDF 可以使用 pandas searchsorted,但需要在执行程序上安装 pandas。 Insted 你可以像这样使用plan python:

    def nearest_holiday(date):
        last_holiday = index.value[0]
        for next_holiday in index.value:
            if next_holiday >= date:
                break
            last_holiday = next_holiday
        if last_holiday > date:
            last_holiday = None
        if next_holiday < date:
            next_holiday = None
        return (last_holiday, next_holiday)
    
    
    from pyspark.sql.types import *
    return_type = StructType([StructField('last_holiday', StringType()), StructField('next_holiday', StringType())])
    
    from pyspark.sql.functions import udf
    nearest_holiday_udf = udf(nearest_holiday, return_type)
    

    并且可以和withColumn一起使用:

    df.withColumn('holiday', nearest_holiday_udf('date')).show(5, False)
    
    +----------+-----------------------+
    |date      |holiday                |
    +----------+-----------------------+
    |2016-01-01|[null,2016-01-03]      |
    |2016-01-02|[null,2016-01-03]      |
    |2016-01-03|[2016-01-03,2016-01-03]|
    |2016-01-04|[2016-01-03,2016-03-03]|
    |2016-01-05|[2016-01-03,2016-03-03]|
    +----------+-----------------------+
    only showing top 5 rows
    

    【讨论】:

    • 谢谢,这看起来很棒。不过,我需要将它移植到 scala ;)
    • 你指的sorted(holidays)操作是什么?它是 pyspark api 吗?
    • 这是python的。它对集合进行排序,因此在 UDF 中我可以遍历它以找到匹配的日期。
    • 我想知道如何获得 2 个单独的列,例如直接从 UDF 展平假期。但是一个UDF只能输出一列?
    • 可以,但是这个列是结构,所以可以在select中简化结构:newdf.select('date', 'holiday.last_holiday', 'holiday.next_holiday')
    猜你喜欢
    • 2021-11-05
    • 1970-01-01
    • 2011-11-12
    • 2018-02-06
    • 2021-12-06
    • 2017-04-07
    • 1970-01-01
    • 2011-06-18
    • 2019-06-30
    相关资源
    最近更新 更多