【问题标题】:How can I add minutes to given timestamp based on my previous row value in pyspark如何根据我在 pyspark 中的前一行值将分钟添加到给定的时间戳
【发布时间】:2019-12-13 14:35:03
【问题描述】:

我有一个 pyspark 数据框

   +----------+----------+---------------------+
   | Activity | Interval |    ReadDateTime     |
   +----------+----------+---------------------+
   |    A     |    1     | 2019-12-13 10:00:00 | 
   |    A     |    2     | 2019-12-13 10:00:00 |
   |    A     |    3     | 2019-12-13 10:00:00 |
   |    B     |    1     | 2019-12-13 11:00:00 | 
   |    B     |    2     | 2019-12-13 11:00:00 |
   |    B     |    3     | 2019-12-13 11:00:00 |
   +--------- +----------+---------------------+

现在我必须根据前一行中的值将 5 分钟添加到 ReadDateTime 列。我预期的数据框如下所示

   +----------+----------+---------------------+
   | Activity | Interval |    ReadDateTime     |
   +----------+----------+---------------------+
   |    A     |    1     | 2019-12-13 10:00:00 | 
   |    A     |    2     | 2019-12-13 10:05:00 |
   |    A     |    3     | 2019-12-13 10:10:00 |
   |    B     |    1     | 2019-12-13 11:00:00 | 
   |    B     |    2     | 2019-12-13 11:05:00 |
   |    B     |    3     | 2019-12-13 11:10:00 |
   +--------- +----------+---------------------+

我不会向对应于间隔 1 的 ReadDateTime 列添加 5 分钟,而我将继续向其他行添加 5 分钟,直到我的活动发生变化

【问题讨论】:

  • 向我们展示您的尝试!
  • 你不能,你需要将数据框提取到对象,更改它们,然后使用覆盖保存到数据框
  • @VladislavVarslavans 我已经找到了解决方案,并在这里发布了

标签: python apache-spark pyspark databricks azure-databricks


【解决方案1】:

感谢 Ali Yesilli 的帖子,我已经找到了解决方案 Adding hours to timestamp in pyspark dynamically.

我首先将我的 ReadDateTime 转换为 unix 时间戳,并且仅当我的 Interval 不等于 1 时才添加 5 分钟。所以我的代码如下所示。

   from pyspark.sql.functions import *

   df = df.withColumn("ReadDateTime1", when(col("Interval") != lit(1),
   col("ReadDateTime") + 
   (col("Interval")*expr("Interval 5 minutes"))).otherwise(col('ReadDateTime')))

【讨论】:

  • 嗨 Saikat,您可以接受它作为答案(单击答案旁边的复选标记,将其从灰色切换为已填充。)。这对其他社区成员可能是有益的。谢谢。
【解决方案2】:

有一个丑陋的方法

from pyspark.sql.functions import *
from pyspark.sql.types import StringType

def update(interval,date):
  if (interval == 1):
    return date
  elif (interval == 2):
    return date + 'add 5 min'
  elif (interval == 3):
    return date + 'add 10 min'

#df.dtypes

my_udf = udf(lambda x,y: update(x,y), StringType())

df.withColumn('updated_realDateTime', my_udf(df.interval, df.realDateTime) ).show(truncate=False)

当然我的更新函数不是你想要的,所以你必须改变它,但它会完成工作(如果模式对于所有间隔都相同,你不需要 elifs,你可以让它动态)

这是为任何有更好答案的人创建数据框的代码

data = [ (1,'2019-12-13 10:00:00'), 
   (2, '2019-12-13 10:00:00'),
   (3, '2019-12-13 10:00:00'),
   (1, '2019-12-13 11:00:00'), 
   (2, '2019-12-13 11:00:00'),
   (3, '2019-12-13 11:00:00')]
df = sqlContext.createDataFrame(data, ['interval','realDateTime']).cache()

【讨论】:

  • 嗨@Benoit,仅供参考,您不需要为此使用udf
  • 嗨@Benoit 我不能使用这个解决方案,因为我有 288 个间隔用于每个活动,我有 10K 的这样的。它会起作用,但不会有效。我发布了一个不使用 UDF 的简单解决方案。请看一下
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-07-02
  • 2020-10-14
  • 2017-10-22
  • 2021-07-11
  • 1970-01-01
  • 2016-09-04
相关资源
最近更新 更多