You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

  1. Convert Raw Data to Parser-Friendly Structure:
    We use map() to turn each AddressRawData record into an AddressData record, setting the unparsed fields (number, road, etc.) to None initially.

  2. Pass Data to the Parser:
    The collect() method pulls the Dataset's records into a local Seq[AddressData], which is exactly what addressParser expects as input.

  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 22:43:12