【问题标题】:Calculate average using Spark Scala使用 Spark Scala 计算平均值
【发布时间】:2017-07-09 14:50:06
【问题描述】:

如何使用以下两个数据集计算 Spark Scala 中每个位置的平均工资?

File1.csv(第4列是工资)

Ram, 30, Engineer, 40000  
Bala, 27, Doctor, 30000  
Hari, 33, Engineer, 50000  
Siva, 35, Doctor, 60000

File2.csv(第 2 列是位置)

Hari, Bangalore  
Ram, Chennai  
Bala, Bangalore  
Siva, Chennai  

以上文件未排序。需要加入这两个文件并找到每个位置的平均工资。我尝试使用以下代码但无法成功。

val salary = sc.textFile("File1.csv").map(e => e.split(","))  
val location = sc.textFile("File2.csv").map(e.split(","))  
val joined = salary.map(e=>(e(0),e(3))).join(location.map(e=>(e(0),e(1)))  
val joinedData = joined.sortByKey()  
val finalData = joinedData.map(v => (v._1,v._2._1._1,v._2._2))  
val aggregatedDF = finalData.map(e=> e.groupby(e(2)).agg(avg(e(1))))    
aggregatedDF.repartition(1).saveAsTextFile("output.txt")  

请帮助提供代码和示例输出的外观。

非常感谢

【问题讨论】:

    标签: scala apache-spark join


    【解决方案1】:

    您可以将 CSV 文件读取为 DataFrame,然后将它们加入并分组以获得平均值:

    val df1 = spark.read.csv("/path/to/file1.csv").toDF(
      "name", "age", "title", "salary"
    )
    
    val df2 = spark.read.csv("/path/to/file2.csv").toDF(
      "name", "location"
    )
    
    import org.apache.spark.sql.functions._
    
    val dfAverage = df1.join(df2, Seq("name")).
      groupBy(df2("location")).agg(avg(df1("salary")).as("average")).
      select("location", "average")
    
    dfAverage.show
    +-----------+-------+
    |   location|average|
    +-----------+-------+
    |Bangalore  |40000.0|
    |  Chennai  |50000.0|
    +-----------+-------+
    

    [更新]用于计算平均尺寸:

    // file1.csv:
    Ram,30,Engineer,40000,600*200
    Bala,27,Doctor,30000,800*400
    Hari,33,Engineer,50000,700*300
    Siva,35,Doctor,60000,600*200
    
    // file2.csv
    Hari,Bangalore
    Ram,Chennai
    Bala,Bangalore
    Siva,Chennai
    
    val df1 = spark.read.csv("/path/to/file1.csv").toDF(
      "name", "age", "title", "salary", "dimensions"
    )
    
    val df2 = spark.read.csv("/path/to/file2.csv").toDF(
      "name", "location"
    )
    
    import org.apache.spark.sql.functions._
    import org.apache.spark.sql.types.IntegerType
    
    val dfAverage = df1.join(df2, Seq("name")).
      groupBy(df2("location")).
      agg(
        avg(split(df1("dimensions"), ("\\*")).getItem(0).cast(IntegerType)).as("avg_length"),
        avg(split(df1("dimensions"), ("\\*")).getItem(1).cast(IntegerType)).as("avg_width")
      ).
      select(
        $"location", $"avg_length", $"avg_width",
        concat($"avg_length", lit("*"), $"avg_width").as("avg_dimensions")
      )
    
    dfAverage.show
    +---------+----------+---------+--------------+
    | location|avg_length|avg_width|avg_dimensions|
    +---------+----------+---------+--------------+
    |Bangalore|     750.0|    350.0|   750.0*350.0|
    |  Chennai|     600.0|    200.0|   600.0*200.0|
    +---------+----------+---------+--------------+
    

    【讨论】:

    • 感谢您的回复。假设该列的尺寸不是工资,而是 600*200(长 * 宽),在这种情况下我如何找到平均值? Ram 600*200 Hari 700*300 等等...
    • @akrockz,请查看扩展答案。
    • 非常感谢@Leo C.. 这就是我要找的.. 最后一个请求.. 目前我的笔记本电脑中没有 Spark 设置.. 你可以发给我吗如果我将输入数据邮寄给您,输出?抱歉问太多了。。谢谢
    • @akrockz,这违反了我的系统使用策略运​​行代码或使用来自未知来源的数据。很抱歉,我无法为您提供帮助。
    • 如果最后我想写入一个 csv 文件,我可以给出以下命令吗? dfAverage.repartition(1).write.csv("output.csv") 这行得通吗?
    【解决方案2】:

    我会使用 DataFrame API,这应该可以:

    val salary = sc.textFile("File1.csv")
                   .map(e => e.split(","))
                   .map{case Seq(name,_,_,salary) => (name,salary)}
                   .toDF("name","salary")
    
    val location = sc.textFile("File2.csv")
                     .map(e => e.split(","))
                     .map{case Seq(name,location) => (name,location)}
                     .toDF("name","location")
    
    import org.apache.spark.sql.functions._
    
    salary
      .join(location,Seq("name"))
      .groupBy($"location")
      .agg(
        avg($"salary").as("avg_salary")
      )
      .repartition(1)
      .write.csv("output.csv")
    

    【讨论】:

    • 所以这里的最终输出如下所示? +------------------------+ |位置 |平均工资 | +------------------------+ |班加罗尔 | 40000 | |钦奈 | 500000 | +------------------------+
    • 还有一个疑问.. 假设该列的尺寸不是工资,而是 600*200(长 * 宽),在这种情况下我如何找到平均值? Ram 600*200 Hari 700*300 等等...
    【解决方案3】:

    我会使用数据框: 首先读取数据帧如:

    val salary = spark.read.option("header", "true").csv("File1.csv")
    val location = spark.read.option("header", "true").csv("File2.csv")
    

    如果您没有标题,则需要将选项设置为“false”并使用 withColumnRenamed 更改默认名称。

    val salary = spark.read.option("header", "false").csv("File1.csv").toDF("name", "age", "job", "salary")
    val location = spark.read.option("header", "false").csv("File2.csv").toDF("name", "location")
    

    现在进行连接:

    val joined = salary.join(location, "name")
    

    最后做平均计算:

    val avg = joined.groupby("location").agg(avg($"salary"))
    

    保存做:

    avg.repartition(1).write.csv("output.csv")
    

    【讨论】:

    • 感谢您的回复。假设该列的尺寸不是工资,而是 600*200(长 * 宽),在这种情况下我如何找到平均值? Ram 600*200 Hari 700*300 等等...
    • 什么意思?您的意思是每个名称多次出现,每个名称有多个列?
    【解决方案4】:

    你可以这样做:

    val salary = sc.textFile("File1.csv").map(_.split(",").map(_.trim))
    val location = sc.textFile("File2.csv").map(_.split(",").map(_.trim))
    val joined = salary.map(e=>(e(0),e(3).toInt)).join(location.map(e=>(e(0),e(1))))
    val locSalary = joined.map(v => (v._2._2, v._2._1))
    val averages = locSalary.aggregateByKey((0,0))((t,e) => (t._1 + 1, t._2 + e),
            (t1,t2) => (t1._1 + t2._1, t1._2 + t2._2)).mapValues(t => t._2/t._1)
    

    然后averages.take(10) 将给出:

    res5: Array[(String, Int)] = Array((Chennai,50000), (Bangalore,40000))
    

    【讨论】:

    • 感谢您的回复。假设该列的尺寸不是工资,而是 600*200(长 * 宽),在这种情况下我如何找到平均值? Ram 600*200 Hari 700*300 等等...
    • 尺寸是否以字符串形式给出?您是要平均面积(长度乘以宽度)还是要为每个维度取平均值?
    • 我想对每个维度进行平均,按位置分组..
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-09-01
    • 2017-08-28
    • 1970-01-01
    • 1970-01-01
    • 2013-01-14
    相关资源
    最近更新 更多