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

在Databricks中统计数组元素频次并转为列的高效方案

高效统计Databricks DataFrame数组列元素频次并转为单独列的方法

针对你遇到的问题,无需逐个统计元素再拼接,推荐使用explode+groupBy+pivot的组合方案,只需少量代码就能完成所有元素的频次统计与列转置,且性能远优于逐个处理的方式。

核心思路

  1. 用explode展开数组列,将每个数组元素转为单独行;
  2. 通过groupBy按names和原items列分组,再用pivot将元素值转成列名,同时统计每个元素的出现次数;
  3. 与原表左连接保留所有行(包括空数组的行),最后用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 08:57:23