【问题标题】:PySpark: counting rows based on current row valuePySpark:根据当前行值计算行数
【发布时间】:2018-06-27 18:19:52
【问题描述】:

我有一个带有“速度”列的 DataFrame。
我能否有效地为每一行添加一列,其中包含 DataFrame 中的行数,以使它们的“速度”在“速度”行的 +/2 范围内?

results = spark.createDataFrame([[1],[2],[3],[4],[5],
                                 [4],[5],[4],[5],[6],
                                 [5],[6],[1],[3],[8],
                                 [2],[5],[6],[10],[12]], 
                                 ['Speed'])

results.show()

+-----+
|Speed|
+-----+
|    1|
|    2|
|    3|
|    4|
|    5|
|    4|
|    5|
|    4|
|    5|
|    6|
|    5|
|    6|
|    1|
|    3|
|    8|
|    2|
|    5|
|    6|
|   10|
|   12|
+-----+

【问题讨论】:

  • 您能添加一个所需输出的样本吗?

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


【解决方案1】:

你可以使用窗口函数:

# Order the window by speed, and look at range [0;+2]
w = Window.orderBy('Speed').rangeBetween(0,2)

# Define a column counting the number of rows containing value Speed+2
results = results.withColumn('count+2',F.count('Speed').over(w)).orderBy('Speed')
results.show()

+-----+-----+
|Speed|count|
+-----+-----+
|    1|    6|
|    1|    6|
|    2|    7|
|    2|    7|
|    3|   10|
|    3|   10|
|    4|   11|
|    4|   11|
|    4|   11|
|    5|    8|
|    5|    8|
|    5|    8|
|    5|    8|
|    5|    8|
|    6|    4|
|    6|    4|
|    6|    4|
|    8|    2|
|   10|    2|
|   12|    1|
+-----+-----+

注意:窗口函数会计算所研究的行本身。您可以通过在计数列中添加 -1 来纠正此问题

results = results.withColumn('count+2',F.count('Speed').over(w)-1).orderBy('Speed')

【讨论】:

  • 非常感谢!我会尝试 :-)。我正在拼命寻找 F.when() 的解决方案,这真的很麻烦。
  • 如果我正在查看的速度范围是十进制,您能给个小费吗?喜欢 +/-0.5 以内的速度?我收到“方法 rangeBetween([class java.lang.Double, class java.lang.Double]) 不存在”错误
  • 确实,看起来pyspark.sql.Window.rangeBetween 只接受整数作为参数。然后,您可以将您的速度列乘以 10,并在 +/-5 范围内工作。 df = df.withColumn("speed_bis" F.col("Speed")*10)
  • 是的,谢谢,这就是我刚刚所做的,它似乎有效。我验证答案。
猜你喜欢
  • 2018-08-07
  • 1970-01-01
  • 2022-07-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多