Spark/Scala禁用spark.sql.function实现客户账户关联数据转CaseClass
实现方案
核心思路
- 全程使用Dataset原生的
groupByKey+mapGroups算子配合Scala标准库集合操作实现,完全不依赖spark.sql.functions下的任何函数 - 左连接后用Option包装账户数据避免空指针异常,按客户ID分组后逐组计算统计指标
- 严格按照业务规则过滤有效账户,不满足条件的账户直接排除,统计指标对应设为0
前置优化建议
因为Scala原生值类型(如Long)不可为空,DataFrame中的null值转成Dataset时会被自动赋值为0,无法区分真实0余额和空值。建议先调整AccountData定义,用Option接收可空字段:
case class AccountData( customerId: String, accountId: String, balance: Option[Long] )
调整后转Dataset的代码不需要改动,Spark会自动把DataFrame的null值映射为None,非null值映射为Some(值)。
完整实现代码
你已有代码之后追加以下逻辑即可:
import scala.collection.mutable.ListBuffer val resultDS = customerAccountsDS // 左连接右表可能为null,用Option包装避免空指针 .map { case (cust, acc) => (cust, Option(acc)) } // 按客户ID分组 .groupByKey(_._1.customerId) // 逐组计算目标输出 .mapGroups { (_, iter) => var customer: CustomerData = null val allAccounts = ListBuffer[AccountData]() // 遍历分组内所有数据,提取客户信息和关联账户 iter.foreach { case (cust, accOpt) => if (customer == null) customer = cust accOpt.foreach(allAccounts.append) } // 过滤有效账户:accountId不为null、balance不为空 val validAccounts = allAccounts.filter(acc => acc.accountId != null && acc.balance.isDefined) if (validAccounts.isEmpty) { CustomerAccountOutput( customerId = customer.customerId, forename = customer.forename, surname = customer.surname, accounts = Seq.empty, numberAccounts = 0, totalBalance = 0L, averageBalance = 0.0 ) } else { val totalBalance = validAccounts.flatMap(_.balance).sum val numberAccounts = validAccounts.size val averageBalance = totalBalance.toDouble / numberAccounts CustomerAccountOutput( customerId = customer.customerId, forename = customer.forename, surname = customer.surname, accounts = validAccounts.toSeq, numberAccounts = numberAccounts, totalBalance = totalBalance, averageBalance = averageBalance ) } }
可选调整说明
如果业务要求无有效账户时totalBalance和averageBalance返回null而非0,只需要调整CustomerAccountOutput对应字段为Option类型,无有效账户时返回None即可。
内容的提问来源于stack exchange,提问作者codeNoob0101
相关产品推荐
相关产品推荐

