Scala Spark基于Seq与Case Class解析地址列报错问题求助
问题根源
报错Can not resolve 'number'的核心原因是:你试图将仅包含addressId、customerId、address三列的addressDF直接转换为AddressData类型的Dataset,但AddressData包含number等原DataFrame不存在的字段,Spark无法完成自动映射,后续操作自然找不到这些字段。
修正方案
1. 修正解析函数的输入类型
原addressParser函数错误地接收Seq[AddressData]作为输入,但原始数据是未解析的AddressRawData,应改为接收AddressRawData并返回AddressData:
def addressParser(rawAddress: AddressRawData): AddressData = { val addressParts = rawAddress.address.split(", ").map(_.trim) // 安全处理拆分后的字段,避免数组越界和类型转换失败 val number = addressParts.headOption.flatMap(_.toIntOption) val road = addressParts.lift(1) val city = addressParts.lift(2) val country = addressParts.lift(3) AddressData( rawAddress.addressId, rawAddress.customerId, rawAddress.address, number, road, city, country ) }
2. 调整数据处理流程
先将原始DataFrame转换为AddressRawData的Dataset,再逐条解析得到AddressData,最后按customerId分组:
// 转换为原始地址数据类型的Dataset(结构完全匹配) val addressRawDS = addressDF.as[AddressRawData] // 解析每条数据得到带拆分字段的AddressData val parsedAddressDS = addressRawDS.map(addressParser) // 按customerId分组,收集用户对应的地址列表 val groupedAddressDF = parsedAddressDS .groupByKey(_.customerId) .mapGroups { case (customerId, addresses) => (customerId, addresses.toSeq) } .toDF("customerId", "address")
3. 完整修正后的代码
object CustomerAddress extends App { val spark = SparkSession.builder().master("local[*]").appName("CustomerAddress").getOrCreate() import spark.implicits._ Logger.getRootLogger.setLevel(Level.WARN) // 定义所有case class case class AddressRawData( addressId: String, customerId: String, address: String ) case class AddressData( addressId: String, customerId: String, address: String, number: Option[Int], road: Option[String], city: Option[String], country: Option[String] ) case class AccountData( customerId: String, accountId: String, balance: Long ) case class CustomerAccountOutput( customerId: String, forename: String, surname: String, accounts: Seq[AccountData] ) // 定义最终输出结构 case class CustomerDocument( customerId: String, forename: String, surname: String, accounts: Seq[AccountData], address: Seq[AddressData] ) // 读取地址CSV,添加quote参数处理带引号的地址字段 val addressDF: DataFrame = spark.read .option("header", "false") .option("quote", "\"") .csv("src/main/resources/address_data.csv") .toDF("addressId", "customerId", "address") val customerAccountDS = spark.read .parquet("src/main/resources/customerAccountOutputDS.parquet") .as[CustomerAccountOutput] // 地址解析函数 def addressParser(rawAddress: AddressRawData): AddressData = { val addressParts = rawAddress.address.split(", ").map(_.trim) val number = addressParts.headOption.flatMap(_.toIntOption) val road = addressParts.lift(1) val city = addressParts.lift(2) val country = addressParts.lift(3) AddressData( rawAddress.addressId, rawAddress.customerId, rawAddress.address, number, road, city, country ) } // 数据处理主流程 val addressRawDS = addressDF.as[AddressRawData] val parsedAddressDS = addressRawDS.map(addressParser) val groupedAddressDF = parsedAddressDS .groupByKey(_.customerId) .mapGroups { case (customerId, addresses) => (customerId, addresses.toSeq) } .toDF("customerId", "address") // 关联数据生成最终结果 val finalDF = customerAccountDS.join(groupedAddressDF, Seq("customerId"), "inner") .select( $"customerId", $"forename", $"surname", $"accounts", $"address" ).as[CustomerDocument] finalDF.show(false) }
关键修改点说明
- 添加
option("quote", "\""):处理CSV中带引号的地址字段,确保拆分时不会把引号内的逗号当作分隔符。 - 用
lift和flatMap(_.toIntOption):安全处理地址拆分后的数组越界和数字转换失败的情况,保证返回合法的Option类型字段值。 - 先解析再分组:先将每条原始地址转换为完整的
AddressData,再按用户ID分组,彻底避免了原代码中类型不匹配的问题。
内容的提问来源于stack exchange,提问作者Mohit Rane
相关产品推荐
相关产品推荐

