【问题标题】:Get start date & end date from the range of timestamp从时间戳范围内获取开始日期和结束日期
【发布时间】:2020-11-19 16:20:33
【问题描述】:

我在 Spark (Scala) 中有一个来自大型 csv 文件的数据框。

Dataframe 是这样的

key| col1 | timestamp            |
---------------------------------
1  | aa  | 2019-01-01 08:02:05.1 |
1  | aa  | 2019-09-02 08:02:05.2 | 
1  | cc  | 2019-12-24 08:02:05.3 |
2  | dd  | 2013-01-22 08:02:05.4 | 

我需要像这样添加两列 start_date 和 end_date

key| col1 | timestamp            | start date              | end date              | 
---------------------------------+---------------------------------------------------
1  | aa  | 2019-01-01 08:02:05.1 | 2017-01-01 08:02:05.1   | 2018-09-02 08:02:05.2 |
1  | aa  | 2019-09-02 08:02:05.2 | 2018-09-02 08:02:05.2   | 2019-12-24 08:02:05.3 |
1  | cc  | 2019-12-24 08:02:05.3 | 2019-12-24 08:02:05.3   | NULL                  |
2  | dd  | 2013-01-22 08:02:05.4 | 2013-01-22 08:02:05.4   | NULL                  |

这里,

对于每一列“键”,end_date 是同一键的下一个时间戳。但是,最新日期的“end_date”应为 NULL。

到目前为止我尝试了什么

我尝试使用窗口函数来计算每个分区的排名

类似的东西

 
  var df = read_csv() 

  //copy timestamp to start_date
  df = df
       .withColumn("start_date", df.col("timestamp"))

  //add null value to the end_date
  df = df.withColumn("end_date", typedLit[Option[String]](None))

  val windowSpec = Window.partitionBy("merge_key_column").orderBy("start_date")
   
  df
  .withColumn("rank", dense_rank()
  .over(windowSpec))
  .withColumn("max", max("rank").over(Window.partitionBy("merge_key_column")))

到目前为止,我还没有得到想要的输出。

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    在这种情况下使用 window lead function

    Example:

    val df=Seq((1,"aa","2019-01-01 08:02:05.1"),(1,"aa","2019-09-02 08:02:05.2"),(1,"cc","2019-12-24 08:02:05.3"),(2,"dd","2013-01-22 08:02:05.4")).toDF("key","col1","timestamp")
    import org.apache.spark.sql.expressions._
    import org.apache.spark.sql.functions._
    import org.apache.spark.sql._
    val df1=df.withColumn("start_date",col("timestamp"))
    val windowSpec = Window.partitionBy("key").orderBy("start_date")
    
    df1.withColumn("end_date",lead(col("start_date"),1).over(windowSpec)).show(10,false)
    //+---+----+---------------------+---------------------+---------------------+
    //|key|col1|timestamp            |start_date           |end_date             |
    //+---+----+---------------------+---------------------+---------------------+
    //|1  |aa  |2019-01-01 08:02:05.1|2019-01-01 08:02:05.1|2019-09-02 08:02:05.2|
    //|1  |aa  |2019-09-02 08:02:05.2|2019-09-02 08:02:05.2|2019-12-24 08:02:05.3|
    //|1  |cc  |2019-12-24 08:02:05.3|2019-12-24 08:02:05.3|null                 |
    //|2  |dd  |2013-01-22 08:02:05.4|2013-01-22 08:02:05.4|null                 |
    //+---+----+---------------------+---------------------+---------------------+
    

    【讨论】:

    • 该解决方案将在col1=aastart_date=2019-12-24 08:02:05.3end_date=null 的位置增加 1 条记录
    猜你喜欢
    • 2020-02-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-03-21
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多