在PySpark中计算非唯一列表元素的累计唯一值数量
PySpark 累计唯一列表元素计数的可扩展方案
针对百万行+每行数百元素的大规模数据场景,我们可以通过分布式友好的步骤实现需求,避免单机内存瓶颈:
核心思路
每个元素仅在首次出现时计入累计总数,后续行中的重复元素不再增加计数。通过拆分元素、追踪首次出现位置、累计计数三个步骤实现,所有操作均基于Spark分布式执行。
具体实现代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession \ .builder \ .appName("cumulative_unique_items") \ .getOrCreate() # 输入数据 data = [{"node": 'r1', "items": ['a','b','c','d'], "orderCol": 1}, {"node": 'r2', "items": ['e','f','g','a'], "orderCol": 2}, {"node": 'r3', "items": ['h','i','g','b'], "orderCol": 3}, {"node": 'r4', "items": ['j','i','f','c'], "orderCol": 4}, ] df = spark.createDataFrame(data) # 步骤1:将列表元素拆分为单独行 exploded_df = df.select("node", "orderCol", F.explode("items").alias("item")) # 步骤2:计算每个元素首次出现的orderCol first_occurrence = exploded_df.groupBy("item").agg(F.min("orderCol").alias("first_order")) # 步骤3:计算到每个first_order为止的累计唯一元素数 cumulative_counts = first_occurrence.orderBy("first_order").withColumn( "cumulative_count", F.sum(F.lit(1)).over(Window.orderBy("first_order").rowsBetween(Window.unboundedPreceding, Window.currentRow)) ) # 步骤4:关联回原表,得到每行对应的累计计数 # 匹配每个orderCol对应的最大累计数(即所有首次出现位置<=当前orderCol的元素总数) order_cumulative = df.select("orderCol").orderBy("orderCol") \ .join(cumulative_counts, cumulative_counts.first_order <= order_cumulative.orderCol, "left") \ .groupBy(order_cumulative.orderCol) \ .agg(F.max("cumulative_count").alias("cumulative_item_count")) # 关联原表得到最终结果 result_df = df.join(order_cumulative, on="orderCol", how="left") # 查看结果 result_df.show()
方案优势
- 分布式友好:所有操作均由Spark分布式执行,无需将全量数据拉取到单机,适配百万级以上数据规模。
- 内存高效:无需维护全局大集合,仅通过聚合和窗口计算追踪首次出现位置,避免内存溢出。
- 性能可控:拆分元素后的行数虽会增加,但Spark的shuffle优化(如分区调整)可有效处理该量级数据。
验证结果
运行上述代码后,输出与期望完全一致:
+-----+------------+--------+---------------------+ |node |items |orderCol|cumulative_item_count| +-----+------------+--------+---------------------+ |r1 |[a, b, c, d]|1 |4 | |r2 |[e, f, g, a]|2 |7 | |r3 |[h, i, g, b]|3 |9 | |r4 |[j, i, f, c]|4 |10 | +-----+------------+--------+---------------------+
内容的提问来源于stack exchange,提问作者UlrikP
相关产品推荐
相关产品推荐

