PySpark统计两列值相同行数新增计数列后去重保留单条记录
PySpark实现方案
实现思路
该需求可以通过两种常用方式实现:
- 窗口函数方案:先按
event+next列分区统计每组行数作为count列,再给同组内的行标记行号,取行号为1的行即可实现每组保留任意一行的要求 - 分组聚合方案:直接按
event和next分组,聚合时同时统计组内行数,以及任选一个组内的id即可,写法更简洁
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("event_next_stat").getOrCreate() # 构造示例数据集 data = [ (1, "A", "X"), (2, "B", "Y"), (3, "C", "Z"), (4, "C", "X"), (5, "A", "X"), (6, "D", "Y"), (7, "B", "Y") ] df = spark.createDataFrame(data, schema=["id", "event", "next"]) # ------------------ 方案1:窗口函数实现 ------------------ # 定义按event、next分区的窗口 window_spec = Window.partitionBy("event", "next") # 新增count列,同时给同组行加随机排序的行号(orderBy用rand保证取任意行) df_with_count = df.withColumn("count", F.count("*").over(window_spec)) \ .withColumn("rn", F.row_number().over(window_spec.orderBy(F.rand()))) # 过滤取每组第一行,删除冗余的rn列 result1 = df_with_count.filter(F.col("rn") == 1).drop("rn") result1.show() # ------------------ 方案2:分组聚合实现(更简洁) ------------------ result2 = df.groupBy("event", "next") \ .agg( F.count("*").alias("count"), # 取组内任意一个id,也可以换成F.min("id")/F.max("id")指定取最小/最大id F.first("id").alias("id") ) \ .select("id", "event", "next", "count") result2.show()
如果需要固定保留每组的最小/最大id,把方案1里的orderBy(F.rand())换成orderBy("id"),或者方案2里的F.first("id")换成F.min("id")/F.max("id")即可,两种方案的输出都和预期结果一致。
内容的提问来源于stack exchange,提问作者lalala
相关产品推荐
相关产品推荐

