【问题标题】:Spark reducer and summation result issue火花减速器和求和结果问题
【发布时间】:2018-01-25 09:42:12
【问题描述】:

这里是示例文件

Department,Designation,costToCompany,State

    Sales,Trainee,12000,UP
    Sales,Lead,32000,AP
    Sales,Lead,32000,LA
    Sales,Lead,32000,TN
    Sales,Lead,32000,AP
    Sales,Lead,32000,TN 
    Sales,Lead,32000,LA
    Sales,Lead,32000,LA
    Marketing,Associate,18000,TN
    Marketing,Associate,18000,TN
    HR,Manager,58000,TN

以 csv 格式生成输出

  • 按部门、职位、州分组

  • 带有 sum(costToCompany) 和 sum(TotalEmployeeCount) 的附加列

结果应该是这样的

Dept,Desg,state,empCount,totalCost
Sales,Lead,AP,2,64000
Sales,Lead,LA,3,96000
Sales,Lead,TN,2,64000

以下是解决方案,写入文件会导致错误。我在这里做错了什么?

第 1 步:加载文件

val file = sc.textFile("data/sales.txt")

步骤#2:创建一个案例类来表示数据

scala> case class emp(Dept:String, Desg:String, totalCost:Double, State:String)
defined class emp

步骤#3:拆分数据并创建emp对象的RDD

scala> val fileSplit = file.map(_.split(","))
scala> val data = fileSplit.map(x => emp(x(0), x(1), x(2).toDouble, x(3)))

第 4 步:将数据转换为 Key/value par with key=(dept, desg,state) 和 value=(1,totalCost)

scala> val keyVals = data.map(x => ((x.Dept,x.Desg,x.State),(1,x.totalCost)))

第 5 步:使用 reduceByKey 进行分组,因为我们还需要员工总数和成本的总和

scala> val results = keyVals.reduceByKey{(a,b) => (a._1+b._1, a._2+b._2)} //(a.count+ b.count, a.cost+b.cost)
results: org.apache.spark.rdd.RDD[((String, String, String), (Int, Double))] = ShuffledRDD[41] at reduceByKey at <console>:55

第 6 步:保存结果

scala> results.repartition(1).saveAsTextFile("data/result")

错误

17/08/16 22:16:59 ERROR executor.Executor: Exception in task 0.0 in stage 20.0 (TID 23)
java.lang.NumberFormatException: For input string: "costToCompany"
    at sun.misc.FloatingDecimal.readJavaFormatString(FloatingDecimal.java:1250)
    at java.lang.Double.parseDouble(Double.java:540)
    at scala.collection.immutable.StringLike$class.toDouble(StringLike.scala:232)
    at scala.collection.immutable.StringOps.toDouble(StringOps.scala:31)
    at $line85.$read$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$anonfun$1.apply(<console>:51)
    at $line85.$read$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$anonfun$1.apply(<console>:51)
    at scala.collection.Iterator$$anon$11.next(Iterator.scala:328)
    at scala.collection.Iterator$$anon$11.next(Iterator.scala:328)
    at org.apache.spark.util.collection.ExternalSorter.insertAll(ExternalSorter.scala:194)
    at org.apache.spark.shuffle.sort.SortShuffleWriter.write(SortShuffleWriter.scala:64)
    at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:73)
    at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:41)
    at org.apache.spark.scheduler.Task.run(Task.scala:89)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:242)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
    at java.lang.Thread.run(Thread.java:745)
17/08/16 22:16:59 WARN scheduler.TaskSetManager: Lost task 0.0 in stage 20.0 (TID 23, localhost, executor driver): java.lang.NumberFormatException: For input string: "costToCompany"

更新 1 忘记删除标题。在这里更新代码。 Save 现在抛出一个不同的错误。另外,需要将标题放回文件中。

scala> val file = sc.textFile("data/sales.txt")
scala> val header = fileSplit.first()
scala> val noHeaderData = fileSplit.filter(_(0) != header(0))
scala> case class emp(Dept:String, Desg:String, totalCost:Double, State:String)
scala> val data = noHeaderData.map(x => emp(x(0), x(1), x(2).toDouble, x(3)))
scala> val keyVals = data.map(x => ((x.Dept,x.Desg,x.State),(1,x.totalCost)))
scala> val resultSpecific = results.map(x => (x._1._1, x._1._2, x._1._3, x._2._1, x._2._2))
scala> resultSpecific.repartition(1).saveASTextFile("data/specific")
<console>:64: error: value saveASTextFile is not a member of org.apache.spark.rdd.RDD[(String, String, String, Int, Double)]
          resultSpecific.repartition(1).saveASTextFile("data/specific")

【问题讨论】:

  • 这部分正确吗? emp(x(0), x(1), x(2).toDouble, x(3)),您使用列表中的前 4 个值,但是,查看您的文件应该是 emp(x(0), x(1), x(4).toDouble, x(2))。另外,您是否从文件中删除了标题?
  • 作业部分正确。我没有删除标题。放置更新 1,现在 saveAsTextFile 有不同的错误。我认为我不需要在保存之前执行另一个地图操作...
  • 结果类似于文件中的(Sales,Lead,AP,2,64000.0)。我怎样才能 1) 添加标题 2) 将条目保存在文件中而不使用 ( 和 ),例如 Sales,Lead,AP,2,64000.0
  • 在这里使用 Spark DataFrames 会更好。
  • 我是大数据概念的新手,来自 C# 和 MVC 背景。我必须深入研究 DataFrames,但现在不是。我需要先整理好我的基本概念。 CCA175 考试也是关于终端的。

标签: scala apache-spark


【解决方案1】:

当您尝试转换为 double 时,costToCompany 字符串不会转换,这就是为什么它在尝试触发动作时卡住的原因。只需从文件中删除第一条记录,然后它就会起作用。您也可以对数据框进行此类操作,这很容易

【讨论】:

    【解决方案2】:

    回答你的问题以及cmets:

    在这种情况下,您可以更轻松地使用数据框,因为您的文件是 csv 格式,您可以使用以下方式加载和保存数据。通过这种方式,您无需担心文件中的行拆分以及标题的处理(加载和保存时)。

    val spark = SparkSession.builder.getOrCreate()
    import spark.implicits._
    
    val df = spark.read
            .format("com.databricks.spark.csv")
            .option("header", "true") //reading the headers
            .load("csv/file/path");
    

    然后,数据框列名称将与文件中的标题相同。而不是reduceByKey(),您可以使用数据框的groupBy()agg()

    val res = df.groupBy($"Department", $"Designation", $"State")
      .agg(count($"costToCompany").alias("empCount"), sum($"costToCompany").alias("totalCost"))
    

    然后保存:

    res.coalesce(1)
      .write.format("com.databricks.spark.csv")
      .option("header", "true")
      .save("results.csv")
    

    【讨论】:

      【解决方案3】:

      错误是直截了当的,它说

      :64: 错误:值 saveASTextFile 不是 org.apache.spark.rdd.RDD[(String, String, String, Int, Double)] resultSpecific.repartition(1).saveASTextFile("data/specific")

      实际上,您没有调用saveASTextFile(...) 的方法,而是saveAsTextFile(???)。您的方法名称有大小写错误。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2017-05-13
        • 1970-01-01
        • 2016-09-01
        • 2015-02-04
        • 1970-01-01
        • 1970-01-01
        • 2021-09-30
        相关资源
        最近更新 更多