如何按PAYMENT_TYPE分组将Spark DataFrame转换为指定JSON格式?
解决Spark DataFrame转指定嵌套JSON格式的方法
核心思路
你之前的实现把字段拼接成了字符串列表,丢失了键值对的结构化信息。需要先保留TRANSACTION_TYPE和TRANSACTION_COUNT的键值对结构,再逐层聚合生成嵌套Map,最后转成目标JSON。
具体实现代码(Scala)
// 1. 按PAYMENT_TYPE分组,将每组内的交易类型和数量转为嵌套Map val groupedDF = res.groupBy($"PAYMENT_TYPE") .agg( map_from_entries(collect_list(struct($"TRANSACTION_TYPE", $"TRANSACTION_COUNT"))).as("transaction_map") ) // 2. 将分组结果转为全局嵌套Map,再生成目标JSON字符串 val targetJson = groupedDF .select(to_json(map_from_entries(collect_list(struct($"PAYMENT_TYPE", $"transaction_map")))).as("final_json")) .first() .getAs[String]("final_json") // 打印结果 println(targetJson)
代码解释
struct($"TRANSACTION_TYPE", $"TRANSACTION_COUNT"):将每一行的交易类型和数量打包成结构化的键值对条目。collect_list(...):收集每个PAYMENT_TYPE分组下的所有交易条目。map_from_entries(...):将条目列表转换为Map结构(比如OFFLINE对应{CHEQUE:5})。map_from_entries(collect_list(struct(...))):将整个分组后的DataFrame转换为一个顶层Map({OFFLINE: {...}, ONLINE: {...}})。to_json(...):将顶层Map直接序列化为你需要的JSON字符串。
输出结果
运行后会得到:
{"OFFLINE": {"CHEQUE": 5}, "ONLINE": {"CREDIT": 135, "DEBIT": 297}}
版本兼容说明
map_from_entries是Spark 2.3及以上版本的内置函数,如果你的Spark版本低于2.3,可以用UDF实现Map合并逻辑:
import org.apache.spark.sql.functions._ import scala.collection.mutable val mergeMaps = udf((maps: Seq[Map[String, Int]]) => { val merged = mutable.Map[String, Int]() maps.foreach(merged ++= _) merged.toMap }) val groupedDF = res.groupBy($"PAYMENT_TYPE") .agg(mergeMaps(collect_list(map($"TRANSACTION_TYPE", $"TRANSACTION_COUNT"))).as("transaction_map")) // 后续生成JSON的步骤同上
内容的提问来源于stack exchange,提问作者SRN
相关产品推荐
相关产品推荐

