Spark Scala中如何对DataFrame进行透视操作(汇率数据场景)
解决方案:Spark Scala实现行转列(透视)操作
你的需求本质是行转列(透视),Spark提供了pivot方法可以高效实现,不需要手动关联两个DataFrame。你之前的代码错误在于:withColumn只能操作当前DataFrame的列,无法直接引用另一个未关联的DataFrame的列,这会触发语法或逻辑错误。
正确实现步骤
- 按货币字段分组,对汇率类型字段执行透视操作,将
01、07转为列 - 聚合取对应汇率值(因为每个货币+类型组合唯一,用
first/max均可) - 重命名列并对数值做两位小数的四舍五入
完整代码示例
import org.apache.spark.sql.functions.{col, first, round} // 假设你的字段常量定义对应示例中的列名 val g_currency_id = "curr" val gf_exchange_rate_applied_type = "type" val gf_exchange_rate_amount = "value" // 1. 分组透视获取原始值 val pivotedDF = balanceConvert .groupBy(g_currency_id) .pivot(gf_exchange_rate_applied_type, Seq("01", "07")) // 指定透视类型值,提升性能 .agg(first(col(gf_exchange_rate_amount))) // 2. 重命名列并处理小数位数 val dfBalance = pivotedDF .withColumnRenamed("01", "t_01") .withColumnRenamed("07", "t_02") .withColumn("t_01", round(col("t_01"), 2)) .withColumn("t_02", round(col("t_02"), 2)) dfBalance.show()
代码说明
pivot方法第二个参数指定要透视的类型值(Seq("01", "07")),避免Spark扫描所有类型值,提升性能first聚合函数用于获取每个货币+类型对应的唯一汇率值,适配你的数据唯一性特点round函数实现数值保留两位小数,匹配你目标DataFrame的格式要求- 全程无需手动创建多个DataFrame再关联,Spark自动处理分组和透视逻辑
输出结果
执行后会得到你期望的DataFrame:
| curr | t_01 | t_02 |
|---|---|---|
| EUR | 0.23 | 0.15 |
| DOL | 0.45 | 0.12 |
内容的提问来源于stack exchange,提问作者FRANCISCO JAVIER ROMERO GARCIA
相关产品推荐
相关产品推荐

