【问题标题】:How to join multiple RDD having Object using scala with case class如何使用带有案例类的scala加入多个具有对象的RDD
【发布时间】:2020-01-14 18:14:15
【问题描述】:

我是 scala 和 spark 的新手,需要以下帮助。
我有三个文件作为地址、邮政编码和大陆,我已经将它们作为 RDD 读取,现在我需要使用 scala 和 spark 找出地址行中具有“stratra”的大陆的数量。

例如:

address             postalcode         continent
stratra 110011     110011 india        india asia
knagar 660011      660011 usa          usa usa
stratra 110012     110012 uk           uk europe
manhatten 669923   669923 usa          usa usa
stratra 220022     220022 srilanka     srilanka asia

所以结果应该是:

((stratra,asia),2)
((stratra,europe),1)

或者如果你能提供更好的选择。

//define three case class in scala:
case class address(line:String,postalcode:Int)
case class postalcode(postalcode:Int,country:String)
case class continent(country:String,continent:String)

val address=sc.tetFile("hdfs://test/address.txt")
val postalcode=sc.tetFile("hdfs://test/postalcode.txt")
val continent=sc.tetFile("hdfs://test/comtinent.txt")

val addressRdd=address.map(x=>x.split(" ")).filetr(v=>v(0)=="stratra").map(line=>Address(line(0),line(1).toInt))

val postRdd=postalcode.map(x=>x.split("")).map(line=>postalcode(line(0).toInt,linr(1))))

val continent=continent.map(x=>x.split("")).map(line=>continent(line(0),linr(1))))

//now I try to join address and postalcode with postalcode
val addresskey=addressRdd.map(line=>(line.postalcode,line))
val postalkey=postalRdd.map(line=>(line.postalcode,line))
val joinaddpostal=addresskey.join(postalkey)

但我没有得到想要的结果来进一步处理
如何做到这一点?提前谢谢

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    你快到了,只是把一些额外的逻辑放在一起。

        // One common case class, so that performing rdd join would be easier
        case class Result(line: String, postalCode: Int, country: String, continent: String)
    
        val address = sc.textFile("D://texts/address.txt")
        val postalCode = sc.textFile("D://texts/postalcode.txt")
        val continent = sc.textFile("D://texts/continent.txt")
    
        // Address keyed by postal code
        val filteredAddress = address.filter(line => line.startsWith("stratra")).map(line => {
          val splits = line.split(" ")
          (splits(1).toInt, Result(splits(0), splits(1).toInt, "", ""))
        })
    
        // Postal codes keyed by postal code
        val postalCodes = postalCode.map(line => {
          val splits = line.split(" ")
          (splits(0).toInt, Result("", splits(0).toInt, splits(1), ""))
        })
    
        // Join address and postal code ON postal code, then convert it to
        // RDD keyed by country
        val addressPostalCode = filteredAddress.join(postalCodes)
                          .map(f => (f._2._2.country, Result(f._2._1.line, f._2._1.postalCode, f._2._2.country, f._2._2.continent)))
    
        // Continent keyed by country
        val continents = continent.map(line => {
          val splits = line.split(" ")
          (splits(0), Result("", 0, splits(0), splits(1)))
        })
    
        // Join continents with address+postalCode rdd
        val merged = continents.join(addressPostalCode)
        val grouped = merged.map(f => ((f._2._2.line, f._2._1.continent), 1)).reduceByKey(_ + _).sortBy(x => x._2, false)
    
        // Check result
        grouped.take(2)
    

    结果是:

    scala> grouped.count
    res39: Long = 2
    
    scala> grouped.foreach(println(_))
    ((stratra,asia),2)
    ((stratra,europe),1)
    

    同样可以通过 Dataframe API 实现,如下所示:

    val address = Seq(("stratra", 110011), ("knagar", 660011), ("stratra", 110012), ("manhatten", 669923), ("stratra", 220022)).toDF("line", "postal")
    
    val postals = Seq(("india", 110011), ("usa", 660011), ("usa", 669923), ("uk", 110012), ("srilanka", 220022)).toDF("country", "postal")
    
    val continent = Seq(("india", "asia"), ("usa", "usa"), ("uk", "europe"), ("usa", "usa"), ("srilanka", "asia")).toDF("country", "continent")
    
    // Filter address starting with "stratra", join with postal on code, join result with continent on country
    val joined = address.filter(row => row.getAs[String]("line").startsWith("stratra")).join(postals, "postal").join(continent, "country")
    
    // Result grouped by line and continent with group count
    val result = joined.select("line", "continent").groupBy("line", "continent").count()
    
    // Check result
    result.show(false)
    

    结果是:

    scala> result.show
    +-------+---------+-----+
    |   line|continent|count|
    +-------+---------+-----+
    |stratra|     asia|    2|
    |stratra|   europe|    1|
    +-------+---------+-----+
    

    【讨论】:

    • 好的,我会测试这个并告诉你结果,谢谢你的负担
    猜你喜欢
    • 2023-02-24
    • 1970-01-01
    • 1970-01-01
    • 2018-10-19
    • 1970-01-01
    • 2023-03-26
    • 1970-01-01
    • 1970-01-01
    • 2016-11-26
    相关资源
    最近更新 更多