【问题标题】:Add new column to dataframe based on previous values and condition根据先前的值和条件向数据框添加新列
【发布时间】:2018-06-01 15:57:23
【问题描述】:

我有示例数据框, 按级别 1 和日期分组后,我得到了结果数据框:

val group_df = qwe.groupBy($"level1",$"date").agg(sum("rel_amount").as("amount"))

+------+----------+------+
|level1|      date|amount|
+------+----------+------+
|     A|2016-03-31|   100|     
|     A|2016-02-28|   100|     
|     A|2016-01-31|   400|     
|     A|2015-12-31|   500|     
|     A|2015-11-30|  1200|     
|     A|2015-10-31|  1300|     
|     A|2014-12-31|   600|     
|     B|2016-03-31|    10|     
|     B|2016-02-28|   300|     
|     B|2016-01-31|   423|     
|     B|2015-12-31|   501|    
|     B|2015-11-30|   234|    
|     B|2015-10-31|  1234|    
|     B|2014-12-31|  3456|    
+------+----------+------+

现在我想添加额外的列(上一列)作为年末,在此列中我需要获取每个组的上一年末金额的值。

例如:对于 level1 :A, date=2016-03-31,该值应为 500,因为它是 2015-12-31 的金额。 同样,对于 date= 2015-12-31,该值应为 600,因为 2014-12-31 的金额。需要计算每一行的上一年年末金额。

预期输出:

+------+----------+------+--------+
|level1|      date|amount|Previous|
+------+----------+------+--------+
|     A|2016-03-31|   100|     500|
|     A|2016-02-28|   100|     500|
|     A|2016-01-31|   400|     500|
|     A|2015-12-31|   500|     600|
|     A|2015-11-30|  1200|     600|
|     A|2015-10-31|  1300|     600|
|     A|2014-12-31|   600|     600|
|     B|2016-03-31|    10|     501|
|     B|2016-02-28|   300|     501|
|     B|2016-01-31|   423|     501|
|     B|2015-12-31|   501|    3456|
|     B|2015-11-30|   234|    3456|
|     B|2015-10-31|  1234|    3456|
|     B|2014-12-31|  3456|    3456|
+------+----------+------+--------+

有人可以帮我解决这个问题吗?

【问题讨论】:

    标签: scala apache-spark spark-dataframe


    【解决方案1】:

    一种方法是使用 UDF 将列 date 操作为 String 以创建一个包含上一个年终值的新列:

    val df = Seq(
      ("A", "2016-03-31", 100),
      ("A", "2016-02-28", 100),
      ("A", "2016-01-31", 400),
      ("A", "2015-12-31", 500),
      ("A", "2015-11-30", 1200),
      ("A", "2015-10-31", 1300),
      ("A", "2014-12-31", 600),
      ("B", "2016-03-31", 10),
      ("B", "2016-02-28", 300),
      ("B", "2016-01-31", 423),
      ("B", "2015-12-31", 501),    
      ("B", "2015-11-30", 234),    
      ("B", "2015-10-31", 1234),   
      ("B", "2014-12-31", 3456)
    ).toDF(
      "level1", "date", "amount"
    )
    
    import org.apache.spark.sql.functions._
    
    def previousEOY = udf( (d: String) => (d.substring(0, 4).toInt - 1).toString + "-12-31" )
    
    val df2 = df.withColumn("previous_eoy", previousEOY($"date"))
    

    为了方便标准 SQL 的标量子查询功能,我将恢复使用 Spark 的 TempView(注意,max() 在子查询中用于满足单行返回):

    df2.createOrReplaceTempView("dfView")
    
    val df3 = spark.sqlContext.sql("""
      SELECT
        level1, date, amount, (
          SELECT max(amount) FROM dfView v2
          WHERE v2.level1 = v1.level1 AND v2.date = v1.previous_eoy
        ) previous
      FROM
        dfView v1
    """)
    
    df3.show
    +------+----------+------+--------+
    |level1|      date|amount|previous|
    +------+----------+------+--------+
    |     A|2016-03-31|   100|     500|
    |     A|2016-02-28|   100|     500|
    |     A|2016-01-31|   400|     500|
    |     A|2015-12-31|   500|     600|
    |     A|2015-11-30|  1200|     600|
    |     A|2015-10-31|  1300|     600|
    |     A|2014-12-31|   600|    null|
    |     B|2016-03-31|    10|     501|
    |     B|2016-02-28|   300|     501|
    |     B|2016-01-31|   423|     501|
    |     B|2015-12-31|   501|    3456|
    |     B|2015-11-30|   234|    3456|
    |     B|2015-10-31|  1234|    3456|
    |     B|2014-12-31|  3456|    null|
    +------+----------+------+--------+
    

    【讨论】:

      【解决方案2】:
      val amount = ss.sparkContext.parallelize(Seq(("B","2014-12-31", 3456))).toDF("level1", "dateY", "amount")
      
      val yearStr = udf((date:String) => {(date.substring(0,4).toInt - 1) +"-12-31" })   
      
      val df3 = amount.withColumn( "p", yearStr($"dateY"))    
      
      df3.show()    
      
      df3.createOrReplaceTempView("dfView")   
      
      val df4 = df3.filter( s => s.getString(1).contains("12-31")).select( $"dateY".as("p"), $"level1",$"amount".as("am"))    
      
      df4.show
      df3.join( df4, Seq("p", "level1"), "left_outer").orderBy("level1", "amount").drop($"p").show()
      

      【讨论】:

      • 感谢为我工作。在最终的数据帧中,它应该是 df4.show df3.join( df4, Seq("p", "level1"), "left_outer").orderBy("level1", "DateY").drop($"p") .show()
      【解决方案3】:

      首先,创建一个年终值的数据框。然后将其加入到您的原始数据框中,其中年份相等。

      【讨论】:

      • 谢谢你的回复,你能解释一下怎么做吗
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-06-20
      • 2018-12-29
      • 2021-03-02
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多