【问题标题】:create a column based on a condition and carrying over previous values根据条件创建列并继承以前的值
【发布时间】:2020-02-23 14:24:33
【问题描述】:

我有以下按“col1”排序的数据框。

+----+----+
|col1|col2|
+----+----+
|   a|   x|
|   a|   x|
|   a|   y|
|   b|   x|
|   b|   z|
|   c|   x|
|   c|   y|
|   d|   z|
|   d|   x|
+----+----+

我想添加一个新列,说“col3”,对于“col1”中唯一组('a','b','c''d')中的每一行,如果“col2”值in ('x' or 'y') 将值加 1,否则如果值为 'z' 或任何其他值结转该值。例如,在“a”的第一行中,因为 col2 是 x,我们通过添加 0 + 1 = 1 来增加 1,在第二行中,因为 col2 再次是 x,我们增加 1 + 1 = 2,依此类推。对于 col1 值为 b(第 4 行)的第二组,我们从新开始,由于 col2 值为 x,我们增加 0 + 1 = 1。在第 5 行,因为 col2 值为 z,我们不增加并取先前的值,即 1 . 在“d”(第8行)的情况下。因为 col2 值不在 x 或 y 中,所以我们不递增并将其保留为 0。

+----+----+----+
|col1|col2|col3|
+----+----+----+
|   a|   x|   1|
|   a|   x|   2|
|   a|   y|   3|
|   b|   x|   1|
|   b|   z|   1|
|   c|   x|   1|
|   c|   y|   2|
|   d|   z|   0|
|   d|   x|   1|
+----+----+----+

无论如何,我可以在不使用 pyspark 中的 UDF 的情况下实现这一目标

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql pyspark-sql pyspark-dataframes


    【解决方案1】:

    使用窗口对col1进行分区,然后用条件表达式新建一列。

    from pyspark.sql.functions import *
    from pyspark.sql import Window
    
    w = Window.partitionBy("col1").rowsBetween(Window.unboundedPreceding, Window.currentRow)
    df.withColumn("col3", sum(when(col("col2").isin("x", "y"), 1).otherwise(0)).over(w)).orderBy("col1").show(10)
    

    代码的结果正是你想要的。

    +----+----+----+
    |col1|col2|col3|
    +----+----+----+
    |   a|   x|   1|
    |   a|   x|   2|
    |   a|   y|   3|
    |   b|   x|   1|
    |   b|   z|   1|
    |   c|   x|   1|
    |   c|   y|   2|
    |   d|   z|   0|
    |   d|   x|   1|
    +----+----+----+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-07-09
      • 1970-01-01
      • 1970-01-01
      • 2020-03-16
      • 1970-01-01
      • 1970-01-01
      • 2020-10-03
      相关资源
      最近更新 更多