【发布时间】: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