Scala中如何对指定Schema的DataFrame统计键值对的分类及总出现次数
解决Scala DataFrame中键值对跨类别统计问题
我来帮你搞定这个统计需求!根据你的DataFrame结构和预期结果,我们可以通过Spark的数组展开、分组聚合和透视操作来实现。下面是完整的步骤和代码:
步骤拆解
- 展开嵌套数组:把
d数组中的每个[k,v]键值对拆成单独一行,让每个键值对都能关联到对应的_id和类别c。 - 统计全局出现ID列表:按每个键值对分组,收集所有出现过的唯一
_id。 - 按类别统计ID列表并透视:按键值对分组后,用
pivot把不同的c值转成独立列,存储对应类别下的ID列表。 - 合并结果并处理空值:把全局统计和类别统计的结果合并,将无数据的类别列转为空数组,保证格式统一。
完整代码实现
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_setvscollect_list:用collect_set避免同一ID重复出现在列表中(比如同一个_id下的同一个键值对重复出现的情况),如果你的数据不会有重复,也可以用collect_list。coalesce处理空值:如果某个键值对在某个类别下没有出现,对应的列会被转为空数组[],而不是null,让结果更规整。
内容的提问来源于stack exchange,提问作者Chaouki
相关产品推荐
相关产品推荐

