【问题标题】:Add days to element inside array in PySpark Dataframe将天数添加到 PySpark Dataframe 中的数组内的元素
【发布时间】:2019-12-07 09:06:43
【问题描述】:

我有一个包含三列的 PySpark 数据框。前两列有数组作为它们的元素,而最后一列给出了最后一列的数组长度。以下是 PySpark 数据框:

+---------------------+---------------------+-----+
|                   c1|                   c2|lenc2|
+---------------------+---------------------+-----+
|[2017-02-14 00:00:00]|[2017-02-24 00:00:00]|    1|
|[2017-01-16 00:00:00]|                   []|    0|
+---------------------+---------------------+-----+

数组包含时间戳数据类型。 lenc2 列表示c1 列中数组的长度。对于lenc2==0 的所有行,c1 列只有一个(时间戳)元素。

对于lenc2==0 的所有行,我想从c1 列中的数组中获取时间戳,将其添加5 天并将其放入c2 行中的数组中。我该怎么做?

这是预期输出的示例:

+---------------------+---------------------+-----+
|                   c1|                   c2|lenc2|
+---------------------+---------------------+-----+
|[2017-02-14 00:00:00]|[2017-02-24 00:00:00]|    1|
|[2017-01-16 00:00:00]|[2017-01-21 00:00:00]|    0|
+---------------------+---------------------+-----+

以下是我到目前为止所尝试的:

df2 = df1.withColumn(
    "c2",
    F.when(F.col("lenc2") == 0, F.array_union(F.col("c1"), F.col("c2"))).otherwise(
        F.col("c2")
    ),
)

【问题讨论】:

    标签: python pyspark pyspark-sql pyspark-dataframes


    【解决方案1】:

    when(…).otherwise(…) 已经正确。

    鉴于您似乎对亚秒精度不感兴趣,您可以将时间戳转换为自 Unix 纪元以来的秒数并添加 5 天的秒数,然后再转换回时间戳:

    from datetime import datetime
    
    from pyspark.sql.functions import *
    
    one_sec_before_leap_time = datetime(2016, 12, 31, 23, 59, 59)
    seconds_in_a_day = 24 * 3600
    
    df = spark.createDataFrame([
        ([one_sec_before_leap_time], [datetime.now()], 1),
        ([one_sec_before_leap_time], [], 0),
    ],
        schema=("c1", "c2", "lenc2"))
    
    
    def add_seconds_to_timestamp(ts_col, seconds_col):
        return to_timestamp(unix_timestamp(ts_col) + seconds_col)
    
    
    df2 = df.withColumn("c2",
                        when(col("lenc2") == 0,
                             array(
                                 add_seconds_to_timestamp(
                                     col("c1").getItem(0),
                                     lit(5 * seconds_in_a_day))))
                        .otherwise(col("c2")))
    df2.show(truncate=False)
    # +---------------------+----------------------------+-----+                      
    # |c1                   |c2                          |lenc2|
    # +---------------------+----------------------------+-----+
    # |[2016-12-31 23:59:59]|[2019-12-07 16:58:32.864176]|1    |
    # |[2016-12-31 23:59:59]|[2017-01-05 23:59:59]       |0    |
    # +---------------------+----------------------------+-----+
    
    

    请注意,当您必须考虑夏令时时,这很可能会给您带来奇怪的结果。最好用 UTC 来表示所有内容,并且只有在输入和输出处才能将 UTC 时间戳正确转换为以本地时区表示的时间。基本上类似于 Unicode 三明治。

    此外,这不考虑leap seconds,如上图所示(2016 年还有一秒钟,使 2016-12-31T12:59:60Z 在技术上有效)。然而,闰秒是出了名的难,因为它没有确切的公式(但谁知道呢,也许有一天我们可以模拟地质和气候事件?)。

    【讨论】:

    • 绝妙的答案。这行得通。感谢您的详细回答和解释。
    猜你喜欢
    • 2022-01-03
    • 1970-01-01
    • 2018-04-29
    • 2019-03-31
    • 2015-05-11
    • 1970-01-01
    • 2017-01-23
    • 1970-01-01
    • 2016-07-31
    相关资源
    最近更新 更多