【发布时间】: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