Spark Dataset提取数组元素并映射至指定CaseClass(禁用spark.sql函数)
解决方案
你当前已经通过joinWith得到了客户和对应账户集合的关联结果,直接基于强类型DataSet的map算子做纯Scala集合运算即可实现需求,全程不需要引入spark.sql.functions包下的任何函数。
前置依赖确认
确保你已经定义好相关CaseClass,且代码开头导入了Spark隐式转换:
import org.apache.spark.sql.SparkSession val spark: SparkSession = SparkSession.builder().getOrCreate() // 必须导入隐式转换,否则DataSet的map操作会缺少编码器报错 import spark.implicits._ // 你已经定义好的相关CaseClass示例(和你实际字段对齐即可) case class Customer(customerId: String, forename: String, surname: String) case class AccountData(customerId: String, accountId: String, balance: Long) case class CustomerAccountOutput( customerId: String, forename: String, surname: String, accounts: Seq[AccountData], numberAccounts: Int, totalBalance: Long, averageBalance: Double )
核心转换代码
直接对你的newTest变量执行map操作,处理空值并计算统计指标即可:
val resultDS = newTest.map { case (customer: Customer, accountTuple: (String, Seq[AccountData])) => // 左连接为空时默认给空账户序列 val accounts = Option(accountTuple).map(_._2).getOrElse(Seq.empty[AccountData]) // 纯Scala集合计算统计值 val accCount = accounts.size val totalBal = accounts.map(_.balance).sum val avgBal = if (accCount == 0) 0.0 else totalBal.toDouble / accCount // 构造输出结构 CustomerAccountOutput( customerId = customer.customerId, forename = customer.forename, surname = customer.surname, accounts = accounts, numberAccounts = accCount, totalBalance = totalBal, averageBalance = avgBal ) } // 输出验证结果 resultDS.show(false)
说明
- 所有统计逻辑都基于原生Scala集合API实现,没有用到任何Spark SQL内置函数,符合培训规则要求
- 用
Option(accountTuple)处理左连接产生的空值,避免空指针异常 - 最终输出的
resultDS就是类型为Dataset[CustomerAccountOutput]的强类型结果,和你要求的输出格式完全一致
内容的提问来源于stack exchange,提问作者codeNoob0101
相关产品推荐
相关产品推荐

