【问题标题】:Creating counter in pyspark在 pyspark 中创建计数器
【发布时间】:2017-03-20 21:23:57
【问题描述】:

如何在 Pyspark 中实现以下 R 代码

l = data.frame(d=c(1,2,4,7,8,15,17,19,20,25,26,29))
l$d2[1]= 0
l$d3[1]=c=1
for(i in 2:nrow(l))
{ l$d2[i]=l$d[i]-l$d[i-1]
  c= ifelse(l$d2[i]<=3,c,c+1)
  l$d3[i]=c
 }
l

如果值大于或等于 3,我想遍历一个列并增加一个计数器。

eg :假设我的列中的元素是

1,2,2,3,2,1,5,2,1

标志应该是: 1,1,1,2,2,2,3,3,3

谢谢

【问题讨论】:

  • 计数器很容易创建,你可以使用spark.apache.org/docs/latest/api/python/…,但是为了帮助你重写代码,你应该先展示你的实验
  • 感谢您的帮助。如果值大于或等于 3,我想遍历一个列并增加一个计数器。例如:假设我的列中的元素是 1,2,2,3,2,1,5,2,1 标志应该是 : 1,1,1,2,2,2,3,3,3
  • 嗨 - 我有类似的问题,你最后是怎么解决的?谢谢!
  • @zhifff : 我已经添加了解决方案

标签: python r apache-spark pyspark


【解决方案1】:

假设以下是输入数据。

输入:

df = spark.createDataFrame([[1,'A',1],[2,'A',2],[3,'A',2],[4,'A',3],[5,'A',2],\
                            [6,'A',5],[7,'B',1],[8,'B',2],[9,'B',5],[10,'B',1]],\
                            ['sl_no','partition','value'])
df.show(10)

  • sl_no - 序列号 [基本上任何定义数据帧顺序的列]
  • partition - 如果需要根据现有列对计数器进行分区,则对列进行分区
  • value - 基于哪个计数器递增的值

输出:

以下代码将为您提供所需的输出。

from pyspark.sql import Window
from pyspark.sql.functions import col, when, sum, lit

threshold= 3

df = df.withColumn("greater",when(col("value")>=lit(threshold),1).otherwise(0))\
       .withColumn("counter",sum("greater").over(Window.partitionBy().orderBy("sl_no")))\
       .withColumn("partitioned_counter",sum("greater").over(Window.partitionBy(["partition"]).orderBy("sl_no")))\
       .orderBy("sl_no")

df.show(10)

  • sl_no - 序列号 [基本上任何定义数据帧顺序的列]
  • partition - 如果需要根据现有列对计数器进行分区,则对列进行分区
  • value - 基于哪个计数器递增的值
  • greater - 检查值是否大于阈值 [在本例中为 3]
  • counter - 当值超过阈值时递增的计数器
  • partitioned_counter - 按分区列分区的计数器

如果您只需要根据列的顺序和阈值创建一个整体计数器,您可以使用上面用于创建 counter 列的代码。

如果用例是为一个分区列/一组分区列单独实现计数器,那么您可以使用用于创建 partitioned_counter 列的代码

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-08-04
    • 1970-01-01
    • 1970-01-01
    • 2016-01-27
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多