Spark中无需使用UDF,通过Map逻辑调用addressParser函数解析地址列生成多列的实现求助
Got it, let's work through this step by step. The main hurdle here is connecting your AddressRawData Dataset to the addressParser function, which expects a Seq[AddressData]. Here's how to bridge that gap and get your address parsing working:
Key Fixes & Steps
First, we need to convert your raw address Dataset into the structure the parser expects, then pass it in, and optionally convert back to a Dataset for Spark operations. I've also added some basic error handling to make the parser more robust (since real-world address data is often messy!).
Full Working Code
object ParseAddress extends App { val spark = SparkSession.builder().master("local[*]").appName("ParseAddress").getOrCreate() import spark.implicits._ case class AddressRawData(addressId: String, address: String) case class AddressData( addressId: String, address: String, number: Option[Int], road: Option[String], city: Option[String], country: Option[String] ) // Updated parser with basic error handling to avoid crashes from malformed addresses def addressParser(unparsedAddresses: Seq[AddressData]): Seq[AddressData] = { unparsedAddresses.map(address => { // Split the address and trim whitespace to handle extra spaces val splitParts = address.address.split(", ").map(_.trim) try { // Remove any leftover quotes from the last field (like in your example) val cleanedCountry = splitParts(3).replaceAll("\"", "") address.copy( number = Some(splitParts(0).toInt), road = Some(splitParts(1)), city = Some(splitParts(2)), country = Some(cleanedCountry) ) } catch { case e: Exception => // Log the issue and return the original record with None values if parsing fails println(s"Warning: Failed to parse address '${address.address}': ${e.getMessage}") address } }) } // Assume you've already loaded your source data into addressDF (e.g., from CSV/DB) // val addressDF = spark.read.format("csv").option("header", true).load("your-data-path.csv") // Step 1: Convert raw data to AddressData structure (initialize optional fields to None) val addressDS: Dataset[AddressRawData] = addressDF.as[AddressRawData] val initialAddressDataDS: Dataset[AddressData] = addressDS.map(rawRecord => AddressData(rawRecord.addressId, rawRecord.address, None, None, None, None) ) // Step 2: Convert Dataset to Seq to match the parser's input type val parsedAddresses: Seq[AddressData] = addressParser(initialAddressDataDS.collect()) // Step 3: (Optional) Convert parsed Seq back to Dataset for distributed Spark operations val parsedAddressDS: Dataset[AddressData] = spark.createDataset(parsedAddresses) // View the final parsed results parsedAddressDS.show(truncate = false) }
Breakdown of Critical Steps
Convert Raw Data to Parser-Friendly Structure:
We usemap()to turn eachAddressRawDatarecord into anAddressDatarecord, setting the unparsed fields (number,road, etc.) toNoneinitially.Pass Data to the Parser:
Thecollect()method pulls the Dataset's records into a localSeq[AddressData], which is exactly whataddressParserexpects as input.Optional: Convert Back to Dataset:
If you need to use Spark's distributed features (like writing to a data lake or running further transformations),spark.createDataset()turns the parsed Seq back into a Dataset.
Quick Notes
- The error handling in the parser prevents your job from crashing if an address is formatted incorrectly (e.g., missing a comma).
- We added
.trim()to split parts and cleaned quotes from the country field to handle the exact format in your example.
内容的提问来源于stack exchange,提问作者Shubh

