You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark统计两列值相同行数新增计数列后去重保留单条记录

PySpark实现方案

实现思路

该需求可以通过两种常用方式实现:

  1. 窗口函数方案:先按event+next列分区统计每组行数作为count列,再给同组内的行标记行号,取行号为1的行即可实现每组保留任意一行的要求
  2. 分组聚合方案:直接按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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.04 03:09:00