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

Scala中如何对指定Schema的DataFrame统计键值对的分类及总出现次数

解决Scala DataFrame中键值对跨类别统计问题

我来帮你搞定这个统计需求!根据你的DataFrame结构和预期结果,我们可以通过Spark的数组展开、分组聚合和透视操作来实现。下面是完整的步骤和代码:

步骤拆解

  1. 展开嵌套数组:把d数组中的每个[k,v]键值对拆成单独一行,让每个键值对都能关联到对应的_id和类别c。
  2. 统计全局出现ID列表:按每个键值对分组,收集所有出现过的唯一_id。
  3. 按类别统计ID列表并透视:按键值对分组后,用pivot把不同的c值转成独立列,存储对应类别下的ID列表。
  4. 合并结果并处理空值:把全局统计和类别统计的结果合并,将无数据的类别列转为空数组,保证格式统一。

完整代码实现

import org.apache.spark.sql.functions.{explode, collect_set, sort_array, concat_ws, pivot, coalesce, array}
import org.apache.spark.sql.types.{ArrayType, LongType}

// 1. 展开数组d,生成键值对、_id、c的中间表
val explodedDf = df.withColumn("key_value_struct", explode($"d"))
  .select(
    $"_id",
    // 生成和示例一致的字符串格式键值对,比如"[k1,v1]"
    concat_ws("", "[", $"key_value_struct.k", ",", $"key_value_struct.v", "]").alias("key_value"),
    $"c"
  )

// 2. 统计每个键值对的全局出现ID列表(排序保证和示例一致)
val totalOccDf = explodedDf.groupBy($"key_value")
  .agg(sort_array(collect_set($"_id")).alias("total_occ"))

// 3. 按类别统计ID列表,并通过pivot转成列
val cOccDf = explodedDf.groupBy($"key_value")
  .pivot("c")
  .agg(sort_array(collect_set($"_id")).alias("occ"))

// 4. 合并结果,将空值列转为空数组
val resultDf = totalOccDf.join(cOccDf, Seq("key_value"), "inner")
  .select(
    $"key_value",
    coalesce($"c_1", array().cast(ArrayType(LongType))).alias("c_1_occ"),
    coalesce($"c_2", array().cast(ArrayType(LongType))).alias("c_2_occ"),
    coalesce($"c_3", array().cast(ArrayType(LongType))).alias("c_3_occ"),
    coalesce($"c_4", array().cast(ArrayType(LongType))).alias("c_4_occ"),
    $"total_occ"
  )

// 查看结果,以[k1,v1]为例
resultDf.filter($"key_value" === "[k1,v1]").show(false)

结果示例

运行上述代码后,你会得到和预期一致的输出:

+---------+--------+--------+--------+--------+---------+
|key_value|c_1_occ |c_2_occ |c_3_occ |c_4_occ |total_occ|
+---------+--------+--------+--------+--------+---------+
|[k1,v1]  |[1,4]   |[2]     |[]      |[]      |[1,2,4]  |
+---------+--------+--------+--------+--------+---------+

关键细节说明

  • sort_array的作用:确保ID列表是升序排列,和你给出的示例格式一致。
  • collect_set vs collect_list:用collect_set避免同一ID重复出现在列表中(比如同一个_id下的同一个键值对重复出现的情况),如果你的数据不会有重复,也可以用collect_list。
  • coalesce处理空值:如果某个键值对在某个类别下没有出现,对应的列会被转为空数组[],而不是null,让结果更规整。

内容的提问来源于stack exchange,提问作者Chaouki

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:00:04