Spark Scala中无法在同一表达式内完成Explode与Select操作
解决方案:在Spark中一次性完成Explode、列重命名与全列选择
我明白你想要在同一个select操作里搞定数组展开、保留所有原有列,同时给展开后的struct字段指定清晰名称并去掉cr:前缀,不想额外写withColumn或者手动逐个列名指定——这完全可以实现,咱直接看具体方案:
核心思路
- 给
explode后的struct列指定一个明确别名,避免Spark默认生成的col名称 - 用
*或自动获取的列列表选择所有原有非数组列,不用手动敲每个列名 - 展开别名后的struct字段,同时批量去掉字段名里的
cr:前缀
完整代码示例
方式一:手动指定非数组列(适合列数少的场景)
import org.apache.spark.sql.functions.{explode, col} val temp = temp1.select( // 选择所有原有非数组列 $"DataPartition", $"TimeStamp", $"TRFCoraxData_instrumentId", $"TRFCoraxData_organizationId", $"FFAction|!|", // 展开Dividends数组并指定别名 explode($"Dividends").alias("dividend") ).select( // 保留原有列 $"DataPartition", $"TimeStamp", $"TRFCoraxData_instrumentId", $"TRFCoraxData_organizationId", $"FFAction|!|", // 展开struct字段,批量替换cr:前缀为空 $"dividend.*".map(field => col(field).alias(field.replace("cr:", ""))): _* ) temp.show(false)
方式二:自动获取非数组列(适合列数多的场景)
如果原有列很多,手动写太麻烦,我们可以自动过滤掉数组列:
import org.apache.spark.sql.functions.{explode, col} // 自动获取所有非Dividends的列 val nonArrayCols = temp1.columns.filter(_ != "Dividends").map(col) val temp = temp1.select( nonArrayCols: _*, explode($"Dividends").alias("dividend") ).select( nonArrayCols: _*, $"dividend.*".map(field => col(field).alias(field.replace("cr:", ""))): _* )
效果说明
这样处理后,你会得到:
- 每个Dividends数组元素展开成单独一行
- 原来带
cr:前缀的字段名,比如cr:AnnouncementDate会变成AnnouncementDate - 所有原有列完整保留,不需要手动逐个指定
输出会类似这样(简化版):
+-----------------+-------------------------+-------------------------+---------------------------+-----------+-------------------------+---------------------+ |DataPartition |TimeStamp |TRFCoraxData_instrumentId|TRFCoraxData_organizationId|FFAction|!||AnnouncementDate |DividendGrossAmount | +-----------------+-------------------------+-------------------------+---------------------------+-----------+-------------------------+---------------------+ |ThirdPartyPrivate|2017-06-07T09:18:33+00:00|8590925624 |4296241518 |O|!| |2009-07-14T00:00:00+00:00|0.35 | |ThirdPartyPrivate|2017-06-07T09:18:33+00:00|8590925624 |4296241518 |O|!| |2008-02-05T00:00:00+00:00|0.3 | +-----------------+-------------------------+-------------------------+---------------------------+-----------+-------------------------+---------------------+
为什么原来的代码会出现col列名?
你之前直接用explode($"Dividends")但没给别名,Spark会默认把这个展开后的struct列命名为col。只要给explode加上.alias("dividend"),就能给它一个明确的名称,后续展开字段也更清晰。
内容的提问来源于stack exchange,提问作者Shailendra
相关产品推荐
相关产品推荐

