【问题标题】:How to apply conditional counts (with reset) to grouped data in PySpark?如何将条件计数(带重置)应用于 PySpark 中的分组数据?
【发布时间】:2019-12-07 23:28:48
【问题描述】:

我有 PySpark 代码,可以有效地将行按数字分组,并在满足特定条件时递增。我无法弄清楚如何有效地将此代码转换为可应用于组的代码。

获取此示例数据框 df

df = sqlContext.createDataFrame(
    [
        (33, [], '2017-01-01'),
        (33, ['apple', 'orange'], '2017-01-02'),
        (33, [], '2017-01-03'),
        (33, ['banana'], '2017-01-04')
    ],
    ('ID', 'X', 'date')
)

这段代码实现了我对这个示例 df 的要求,即按日期排序并创建在 size 列返回 0 时递增的组 ('grp')。

df \
.withColumn('size', size(col('X'))) \
.withColumn(
    "grp", 
    sum((col('size') == 0).cast("int")).over(Window.orderBy('date'))
).show()

这部分基于Pyspark - Cumulative sum with reset condition

现在我要做的是将相同的方法应用于具有多个 ID 的数据框 - 实现看起来像的结果

df2 = sqlContext.createDataFrame(
    [
        (33, [], '2017-01-01', 0, 1),
        (33, ['apple', 'orange'], '2017-01-02', 2, 1),
        (33, [], '2017-01-03', 0, 2),
        (33, ['banana'], '2017-01-04', 1, 2),
        (55, ['coffee'], '2017-01-01', 1, 1),
        (55, [], '2017-01-03', 0, 2)
    ],
    ('ID', 'X', 'date', 'size', 'group')
)

为清晰起见进行编辑

1) 对于每个 ID 的第一个日期 - 组应该是 1 - 无论在任何其他列中显示什么。

2) 但是,对于每个后续日期,我需要检查大小列。如果大小列是 0,那么我增加组号。如果它是任何非零的正整数,那么我继续之前的组号。

我在 pandas 中看到了一些处理此问题的方法,但我很难理解 pyspark 中的应用程序以及 pandas 与 spark 中分组数据不同的方式(例如,我是否需要使用称为 UADF 的东西?)

【问题讨论】:

  • 为什么咖啡行的组值为1?不应该是0吗?
  • @cronoik 它是 1,因为 ID 已更改,这一行是该 ID 的第一个日期 - 所以它应该是第 1 组。
  • 所以起始值是X的大小,每次X的大小为0时递增?
  • 起始值始终为 1 - 第一个 ID + 日期。然后该值仅在 X 的大小为 0 时增加。第一个 ID + 日期的 X 的大小无关紧要(在我的实际数据中,它始终为 0 或缺失,但我试图在这里简化) .我会将这些信息编辑到主要问题中

标签: pyspark pyspark-sql


【解决方案1】:

通过检查size 是否为零或该行是第一行来创建列zero_or_first。然后sum

df2 = sqlContext.createDataFrame(
    [
        (33, [], '2017-01-01', 0, 1),
        (33, ['apple', 'orange'], '2017-01-02', 2, 1),
        (33, [], '2017-01-03', 0, 2),
        (33, ['banana'], '2017-01-04', 1, 2),
        (55, ['coffee'], '2017-01-01', 1, 1),
        (55, [], '2017-01-03', 0, 2),
        (55, ['banana'], '2017-01-01', 1, 1)
    ],
    ('ID', 'X', 'date', 'size', 'group')
)


w = Window.partitionBy('ID').orderBy('date')
df2 = df2.withColumn('row', F.row_number().over(w))
df2 = df2.withColumn('zero_or_first', F.when((F.col('size')==0)|(F.col('row')==1), 1).otherwise(0))
df2 = df2.withColumn('grp', F.sum('zero_or_first').over(w))
df2.orderBy('ID').show()

这里是输出。您可以看到该列group == grp。其中group 是预期结果。

+---+---------------+----------+----+-----+---+-------------+---+
| ID|              X|      date|size|group|row|zero_or_first|grp|
+---+---------------+----------+----+-----+---+-------------+---+
| 33|             []|2017-01-01|   0|    1|  1|            1|  1|
| 33|       [banana]|2017-01-04|   1|    2|  4|            0|  2|
| 33|[apple, orange]|2017-01-02|   2|    1|  2|            0|  1|
| 33|             []|2017-01-03|   0|    2|  3|            1|  2|
| 55|       [coffee]|2017-01-01|   1|    1|  1|            1|  1|
| 55|       [banana]|2017-01-01|   1|    1|  2|            0|  1|
| 55|             []|2017-01-03|   0|    2|  3|            1|  2|
+---+---------------+----------+----+-----+---+-------------+---+


【讨论】:

    【解决方案2】:

    我添加了一个窗口函数,并在每个 ID 中创建了一个索引。然后我扩展了条件语句以引用该索引。以下似乎产生了我想要的输出数据框 - 但我很想知道是否有更有效的方法来做到这一点。

    window = Window.partitionBy('ID').orderBy('date')
    df \
    .withColumn('size', size(col('X'))) \
    .withColumn('index', rank().over(window).alias('index')) \
    .withColumn(
        "grp", 
        sum(((col('size') == 0) | (col('index') == 1)).cast("int")).over(window)
    ).show()
    

    产生

    +---+---------------+----------+----+-----+---+
    | ID|              X|      date|size|index|grp|
    +---+---------------+----------+----+-----+---+
    | 33|             []|2017-01-01|   0|    1|  1|
    | 33|[apple, orange]|2017-01-02|   2|    2|  1|
    | 33|             []|2017-01-03|   0|    3|  2|
    | 33|       [banana]|2017-01-04|   1|    4|  2|
    | 55|       [coffee]|2017-01-01|   1|    1|  1|
    | 55|             []|2017-01-03|   0|    2|  2|
    +---+---------------+----------+----+-----+---+
    

    【讨论】:

    • 你定义的窗口是什么?
    • 加上(55, ['banana'], '2017-01-01', 1, 1)这行不行。问题是我认为应该是(55, ['coffee'], '2017-01-01', 1, 0)
    • @cronoik 将 Window 添加到我的答案中
    猜你喜欢
    • 1970-01-01
    • 2019-12-23
    • 2020-08-01
    • 2023-04-07
    • 2019-10-12
    • 1970-01-01
    • 1970-01-01
    • 2019-06-18
    • 1970-01-01
    相关资源
    最近更新 更多