【问题标题】:pyspark how to return the average of a column based on the value of another column?pyspark如何根据另一列的值返回一列的平均值?
【发布时间】:2020-05-25 03:13:54
【问题描述】:

我不认为这会很困难,但我无法理解如何在我的 spark 数据框中取列的平均值。

数据框如下所示:

+-------+------------+--------+------------------+
|Private|Applications|Accepted|              Rate|
+-------+------------+--------+------------------+
|    Yes|         417|     349|0.8369304556354916|
|    Yes|        1899|    1720|0.9057398630858347|
|    Yes|        1732|    1425|0.8227482678983834|
|    Yes|         494|     313|0.6336032388663968|
|     No|        3540|    2001|0.5652542372881356|
|     No|        7313|    4664|0.6377683577191303|
|    Yes|         619|     516|0.8336025848142165|
|    Yes|         662|     513|0.7749244712990937|
|    Yes|         761|     725|0.9526938239159002|
|    Yes|        1690|    1366| 0.808284023668639|
|    Yes|        6075|    5349|0.8804938271604938|
|    Yes|         632|     494|0.7816455696202531|
|     No|        1208|     877|0.7259933774834437|
|    Yes|       20192|   13007|0.6441660063391442|
|    Yes|        1436|    1228|0.8551532033426184|
|    Yes|         392|     351|0.8954081632653061|
|    Yes|       12586|    3239|0.2573494358811378|
|    Yes|        1011|     604|0.5974282888229476|
|    Yes|         848|     587|0.6922169811320755|
|    Yes|        8728|    5201|0.5958982584784601|
+-------+------------+--------+------------------+

当Private 等于“是”时,我想返回Rate 列的平均值。我该怎么做?

【问题讨论】:

  • 过滤然后计算平均值?
  • 没错,我不知道如何计算平均值。

标签: python dataframe apache-spark pyspark


【解决方案1】:

执行相同操作的第三个版本是:

from pyspark.sql.functions import col, avg
df_avg = df.filter(df["Private"] == "Yes").agg(avg(col("Rate")))
df_avg.show()

【讨论】:

  • 我认为这与 Vishnudev 的答案一起工作,但答案显示不正确。就像我之前提到的,打印这个给我:DataFrame[avg(Rate): double],我需要查看号码。我很好奇你提到的显示功能,看来我需要先导入一些东西才能工作。
  • df_avg.show() 给了我一个很长的错误(我在尝试非常不同的事情时经常看到这个错误)
【解决方案2】:

这将在 scala 中工作。 pyspark 代码应该非常相似。

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

val df = List(
("yes", 10),
("yes", 30),
("No", 40)).toDF("private", "rate")

val df = l.toDF(List("private", "rate"))

val window =Window.partitionBy($"private")

df.
    withColumn("avg", 
                when($"private" === "No", null).
                otherwise(avg($"rate").over(window))
            ).
    show()

输入 DF

+-------+----+
|private|rate|
+-------+----+
|    yes|  10|
|    yes|  30|
|     No|  40|
+-------+----+

输出df

+-------+----+----+
|private|rate| avg|
+-------+----+----+
|     No|  40|null|
|    yes|  10|20.0|
|    yes|  30|20.0|
+-------+----+----+

【讨论】:

    【解决方案3】:

    试试

    df.filter(df['Private'] == 'Yes').agg({'Rate': 'avg'}).collect()[0]
    

    【讨论】:

    • 我不确定过滤器是否以这种方式工作。过滤器适用于索引和列标题。
    • 我试过privateRate = df.filter(df['Private'] == 'Yes').agg({'Rate': 'avg'}),但print(privateRate)返回DataFrame[avg(Rate): double]。似乎很接近,但我需要查看实际数字@Vishnudev
    • spyder 警告我“未定义的名称 display”
    【解决方案4】:

    试试:

    from pyspark.sql.functions import col, mean, lit
    
    df.where(col("Private")==lit("Yes")).select(mean(col("Rate"))).collect()
    

    【讨论】:

    • spyder 给我一个代码分析警告,“'lit' 可能未定义或从星形导入中定义”。这也会导致运行代码时出错
    • 欧,现在试试 - F.lit(...) 应该这样做
    • 我收到一个很长的错误。我在这个项目上见过几次,但并不总是出现。我开始怀疑问题是否不在于我的代码的这一部分,而是其他问题?这会很奇怪,因为我可以毫无问题地跑到这些线。
    • 我开始做事了 - from pyspark.sql import functions as F 是这里的问题(一般情况下不应该是这种情况)。您现在可以检查 - 但如果仍然不是,那是一些环境差异,我无法重现......
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-08-21
    • 2020-10-27
    • 1970-01-01
    • 2020-05-13
    • 1970-01-01
    • 1970-01-01
    • 2019-05-23
    相关资源
    最近更新 更多