【问题标题】:How to get value from previous group in spark?如何从火花中的前一组获得价值?
【发布时间】:2021-08-19 02:12:24
【问题描述】:

我需要在 spark 中获取前一个组的值并将其设置为当前组。 我怎样才能做到这一点? 我必须按计数而不是 TEXT_NUM 排序。

无法按 TEXT_NUM 排序,因为事件会按时间重复,如计数 10 和 11 所示。

我正在尝试使用以下代码:

   val spark = SparkSession.builder()
      .master("spark://spark-master:7077")
      .getOrCreate()

    val df = spark
      .createDataFrame(
        Seq[(Int, String, Int)](
          (0, "", 0),
          (1, "", 0),
          (2, "A", 1),
          (3, "A", 1),
          (4, "A", 1),
          (5, "B", 2),
          (6, "B", 2),
          (7, "B", 2),
          (8, "C", 3),
          (9, "C", 3),
          (10, "A", 1),
          (11, "A", 1)
        ))
      .toDF("count", "TEXT", "TEXT_NUM")

    val w1 = Window
      .orderBy("count")
      .rangeBetween(Window.unboundedPreceding, -1)
    df
      .withColumn("LAST_VALUE", last("TEXT_NUM").over(w1))
      .orderBy("count")
      .show()

结果:

+-----+----+--------+----------+
|count|TEXT|TEXT_NUM|LAST_VALUE|
+-----+----+--------+----------+
|    0|    |       0|      null|
|    1|    |       0|         0|
|    2|   A|       1|         0|
|    3|   A|       1|         1|
|    4|   A|       1|         1|
|    5|   B|       2|         1|
|    6|   B|       2|         2|
|    7|   B|       2|         2|
|    8|   C|       3|         2|
|    9|   C|       3|         3|
|   10|   A|       1|         3|
|   11|   A|       1|         1|
+-----+----+--------+----------+

想要的结果:

+-----+----+--------+----------+
|count|TEXT|TEXT_NUM|LAST_VALUE|
+-----+----+--------+----------+
|    0|    |       0|      null|
|    1|    |       0|      null|
|    2|   A|       1|         0|
|    3|   A|       1|         0|
|    4|   A|       1|         0|
|    5|   B|       2|         1|
|    6|   B|       2|         1|
|    7|   B|       2|         1|
|    8|   C|       3|         2|
|    9|   C|       3|         2|
|   10|   A|       1|         3|
|   11|   A|       1|         3|
+-----+----+--------+----------+

【问题讨论】:

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


    【解决方案1】:

    考虑使用 Window 函数 last(columnName, ignoreNulls) 回填 nulls 在组边界处包含先前“text_num”的列中,如下所示:

    val df = Seq(
      (0, "", 0), (1, "", 0),
      (2, "A", 1), (3, "A", 1), (4, "A", 1),
      (5, "B", 2), (6, "B", 2), (7, "B", 2),
      (8, "C", 3), (9, "C", 3),
      (10, "A", 1), (11, "A", 1)
    ).toDF("count", "text", "text_num")
    
    import org.apache.spark.sql.expressions.Window
    val w1 = Window.orderBy("count")
    val w2 = w1.rowsBetween(Window.unboundedPreceding, 0)
    
    df.
      withColumn("prev_num", lag("text_num", 1).over(w1)).
      withColumn("last_change", when($"text_num" =!= $"prev_num", $"prev_num")).
      withColumn("last_value", last("last_change", ignoreNulls=true).over(w2)).
      show
    /*
    +-----+----+--------+--------+-----------+----------+
    |count|text|text_num|prev_num|last_change|last_value|
    +-----+----+--------+--------+-----------+----------+
    |    0|    |       0|    null|       null|      null|
    |    1|    |       0|       0|       null|      null|
    |    2|   A|       1|       0|          0|         0|
    |    3|   A|       1|       1|       null|         0|
    |    4|   A|       1|       1|       null|         0|
    |    5|   B|       2|       1|          1|         1|
    |    6|   B|       2|       2|       null|         1|
    |    7|   B|       2|       2|       null|         1|
    |    8|   C|       3|       2|          2|         2|
    |    9|   C|       3|       3|       null|         2|
    |   10|   A|       1|       3|          3|         3|
    |   11|   A|       1|       1|       null|         3|
    +-----+----+--------+--------+-----------+----------+
    */
    

    中间列保留在输出中以供参考。如果不需要它们,请丢弃它们。

    【讨论】:

      猜你喜欢
      • 2016-03-12
      • 1970-01-01
      • 2019-08-01
      • 1970-01-01
      • 2016-12-26
      • 2019-12-28
      • 2021-07-28
      • 1970-01-01
      • 2013-11-10
      相关资源
      最近更新 更多