【问题标题】:PySpark Window Function ComprehensionPySpark 窗口函数理解
【发布时间】:2018-03-14 18:30:41
【问题描述】:

我一直在使用 PySpark 在一个相当大的 RDD 中运行大量计算,每个块如下所示:

ID  CHK  C1 Flag1   V1  V2  C2  Flag2   V3  V4  
341 10  100 TRUE    10  10  150 FALSE   10  14
341 9   100 TRUE    10  10  150 FALSE   10  14
341 8   100 TRUE    14  14  150 FALSE   10  14
341 7   100 TRUE    14  14  150 FALSE   10  14
341 6   100 TRUE    14  14  150 FALSE   10  14
341 5   100 TRUE    14  14  150 FALSE   10  14
341 4   100 TRUE    14  14  150 FALSE   12  14
341 3   100 TRUE    14  14  150 FALSE   14  14
341 2   100 TRUE    14  14  150 FALSE   14  14
341 1   100 TRUE    14  14  150 FALSE   14  14
341 0   100 TRUE    14  14  150 FALSE   14  14

我有很多次出现的 ID(它取决于 C1 值,例如从 100 到 130 等等,对于许多 C1,对于每个整数,我有一组 11 行,就像上面的那样)并且我有很多 ID .我需要做的是在每行的组中应用一个公式并添加两列来计算:

D1 = ((row.V1 - prev_row.V1)/2)/((row.V2 + prev_row.V2)/2)
D2 = ((row.V3 - prev_row.V3)/2)/((row.V4 + prev_row.V4)/2)

我所做的(正如我在这篇有用的文章中发现的:https://arundhaj.com/blog/calculate-difference-with-previous-row-in-pyspark.html)是定义一个窗口:

my_window = Window.partitionBy().orderBy(desc("CHK"))

对于每个中间计算,我创建了一个“temp”列:

df = df.withColumn("prev_V1", lag(df.V1).over(my_window))
df = df.withColumn("prev_V21", lag(df.TA1).over(my_window))
df = df.withColumn("prev_V3", lag(df.SSQ2).over(my_window))
df = df.withColumn("prev_V4", lag(df.TA2).over(my_window))
df = df.withColumn("Sub_V1", F.when(F.isnull(df.V1 - df.prev_V1), 0).otherwise((df.V1 - df.prev_V1)/2))
df = df.withColumn("Sub_V2", (df.V2 + df.prev_V2)/2)
df = df.withColumn("Sub_V3", F.when(F.isnull(df.V3 - df.prev_V3), 0).otherwise((df.V3 - df.prev_V3)/2))
df = df.withColumn("Sub_V4", (df.V4 + df.prev_V4)/2)
df = df.withColumn("D1", F.when(F.isnull(df.Sub_V1 / df.Sub_V2), 0).otherwise(df.Sub_V1 / df.Sub_V2))
df = df.withColumn("D2", F.when(F.isnull(df.Sub_V3 / df.Sub_V4), 0).otherwise(df.Sub_V3 / df.Sub_V4))

最后我去掉了临时列:

final_df = df.select(*columns_needed)

花了很长时间,我不断得到:

WARN WindowExec: No Partition Defined for Window operation! Moving all data to a single partition, this can cause serious performance degradation.

我知道我没有正确执行此操作,因为上面的代码块位于几个 for 循环中,以便对所有 ID 进行计算,即循环使用:

unique_IDs = list(df1.toPandas()['ID'].unique())

但是在研究了 PySpark Window 函数的更多信息后,我相信通过正确设置窗口 partitionBy(),我可以更轻松地获得相同的结果。

我查看了Avoid performance impact of a single partition mode in Spark window functions,但我仍然不确定如何正确设置窗口分区以使其正常工作。

有人可以就我如何解决这个问题向我提供一些帮助或见解吗?

谢谢

【问题讨论】:

  • 我需要做的是在每一行的组中应用一个公式——你的意思是 ID 和 C1 组?如果是这样,那么您的分区应该在这些列上。

标签: python function pyspark spark-dataframe


【解决方案1】:

我假设必须将公式应用于每组 ID(这就是我选择在“ID”上进行分区的原因)。

您可以避免使用类似这样的“临时”列:

# used to define the lag of a specific column
w_lag=Window.partitionBy("id","C1").orderBy(desc('chk'))


df = df.withColumn('D1',((df.V1-F.lag(df.V1).over(w_lag))/2)\
                       /((df.V2+F.lag(df.V2).over(w_lag))/2))

df = df.withColumn('D2',((df.V3-F.lag(df.V3).over(w_lag))/2)\
                       /((df.V4+F.lag(df.V4).over(w_lag))/2))

结果是:

+---+---+---+-----+---+---+---+-----+---+---+-------------------+-------------------+
| id|chk| C1|Flag1| V1| V2| C2|Flag2| V3| V4|                 D1|                 D2|
+---+---+---+-----+---+---+---+-----+---+---+-------------------+-------------------+
|341| 10|100| true| 10| 10|150| true| 10| 14|               null|               null|
|341|  9|100| true| 10| 10|150| true| 10| 14|                0.0|                0.0|
|341|  8|100| true| 14| 14|150| true| 10| 14|0.16666666666666666|                0.0|
|341|  7|100| true| 14| 14|150| true| 10| 14|                0.0|                0.0|
|341|  6|100| true| 14| 14|150| true| 10| 14|                0.0|                0.0|
|341|  5|100| true| 14| 14|150| true| 10| 14|                0.0|                0.0|
|341|  4|100| true| 14| 14|150| true| 12| 14|                0.0|0.07142857142857142|
|341|  3|100| true| 14| 14|150| true| 14| 14|                0.0|0.07142857142857142|
|341|  2|100| true| 14| 14|150| true| 14| 14|                0.0|                0.0|
|341|  1|100| true| 14| 14|150| true| 14| 14|                0.0|                0.0|
|341|  0|100| true| 14| 14|150| true| 14| 14|                0.0|                0.0|
+---+---+---+-----+---+---+---+-----+---+---+-------------------+-------------------+

【讨论】:

  • 这似乎接近我的想法,我唯一担心的是这是否适用于整个数据集,例如下一个“数据块”也将具有 id = 341 和 chk = 10 -> 0 但 C1 将是 110。你认为这样做 w_lag=Window.partitionBy("id","C1").orderBy(desc( 'chk')) 会解决这个问题吗?最后,我认为当 D1 = 0.16666 时应该是:((14-10)/2)/(14+14/2) = 2/14 = 0.1428
  • 我刚刚编辑了我的答案。确实: Window.partitionBy("id","C1") 应该更适合您想要的。至于 D1 = 0.16666 :我认为是 ((14-10)/2)/((14+10)/2)=2/12=0.1666
  • 你是对的,我的错!非常感谢您的回复!
猜你喜欢
  • 2022-01-25
  • 2019-08-16
  • 2019-09-21
  • 2023-02-23
  • 2017-03-17
  • 2016-05-10
  • 2018-01-25
  • 2018-07-11
  • 1970-01-01
相关资源
最近更新 更多