Spark Scala如何将map<string,string>列转为_id结尾的独立列
实现方案
你可以通过以下两种方式实现需求,可根据你的实际场景选择:
方案1:动态Pivot(适用ATTRIBUTES中_id结尾的键不固定的场景)
不需要提前知道具体键名,Spark会自动识别所有符合规则的键生成列:
import org.apache.spark.sql.functions._ // 展开map并过滤符合后缀要求的键值对 val filteredExplodedDf = batch_df .select($"event_date", explode($"ATTRIBUTES")) .filter(lower($"key").endsWith("_id")) // 兼容大小写后缀,如Account_Id、mem_id都可以匹配 // 按event_date分组后转宽表 val resultDf = filteredExplodedDf .groupBy("event_date") .pivot("key") .agg(first($"value")) // 同一event_date下同一个key只有一个值时用first,多个值可根据需求换collect_list/collect_set等聚合函数 resultDf.show()
方案2:直接提取列(适用键已知的场景,性能更高)
如果提前知道所有要提取的_id结尾的键名,无需shuffle操作,性能远高于Pivot方案:
import org.apache.spark.sql.functions._ // 提前定义所有目标键名 val targetIdKeys = Seq("SYST_id", "RECVR_id", "Account_Id", "Vb_id", "SYS_INFO_id", "mem_id") // 构造查询列,直接从map中取对应键的值 val selectExpr = col("event_date") +: targetIdKeys.map(key => $"ATTRIBUTES".getItem(key).alias(key)) val resultDf = batch_df.select(selectExpr: _*) resultDf.show()
内容的提问来源于stack exchange,提问作者SanjanaSanju
相关产品推荐
相关产品推荐

