在Databricks中统计数组元素频次并转为列的高效方案
高效统计Databricks DataFrame数组列元素频次并转为单独列的方法
针对你遇到的问题,无需逐个统计元素再拼接,推荐使用explode+groupBy+pivot的组合方案,只需少量代码就能完成所有元素的频次统计与列转置,且性能远优于逐个处理的方式。
核心思路
- 用
explode展开数组列,将每个数组元素转为单独行; - 通过
groupBy按names和原items列分组,再用pivot将元素值转成列名,同时统计每个元素的出现次数; - 与原表左连接保留所有行(包括空数组的行),最后用
fillna将空值填充为0。
完整实现代码
from pyspark.sql import functions as F # 假设你的dummy_df已创建,这里给出示例数据初始化代码(可替换为你的实际数据) data = [ ("Ash", ["c1","c2","c2","c3"]), ("Bob", ["c1","c2"]), ("May", []), ("Amy", ["c2","c3","c3"]) ] dummy_df = spark.createDataFrame(data, ["names", "items"]) # 步骤1:展开数组并转置统计列 # 若提前知道所有可能的item(如c1、c2、c3),可以在pivot时传入列表,提升大数据量下的性能 target_items = ["c1", "c2", "c3"] count_pivot_df = dummy_df.select("names", "items", F.explode("items").alias("item")) \ .groupBy("names", "items") \ .pivot("item", target_items) \ .agg(F.count("item")) \ # 重命名列名为xxx_count格式 .withColumns({f"{item}_count": F.col(item) for item in target_items}) \ .drop(*target_items) # 步骤2:关联原表并填充空值为0 result_df = dummy_df.join(count_pivot_df, on=["names", "items"], how="left") \ .fillna(0, subset=[f"{item}_count" for item in target_items]) # 调整列顺序(可选,按你的期望输出排列) result_df = result_df.select("names", "items", "c1_count", "c2_count", "c3_count") # 查看结果 result_df.show()
关键优势
- 性能高效:仅需一次
explode和分组统计操作,避免了多次统计与拼接的冗余计算,即使元素数量达到50-100个也能轻松处理; - 扩展性强:只需修改
target_items列表即可适配不同的元素集合,无需修改核心逻辑; - 覆盖边界场景:通过左连接和
fillna确保空数组的行所有统计列都为0,符合你的输出要求。
内容的提问来源于stack exchange,提问作者Joe
相关产品推荐
相关产品推荐

