【发布时间】:2017-02-26 03:02:15
【问题描述】:
我正在运行一个简单的 sparkSQL 查询,它在 2 个数据集上进行匹配,每个数据集大约 500GB。所以整个数据大约是 1TB。
val adreqPerDeviceid = sqlContext.sql("select count(Distinct a.DeviceId) as MatchCount from adreqdata1 a inner join adreqdata2 b ON a.DeviceId=b.DeviceId ")
adreqPerDeviceid.cache()
adreqPerDeviceid.show()
在加载数据(分配 10k 个任务)之前,作业工作正常。
200 个任务分配在.cache 行。它失败的地方!我知道我没有缓存大量数据,它只是一个数字,为什么它会在这里失败。
以下是错误详情:
在 org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1283) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1271) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1270) 在 scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59) 在 scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:47) 在 org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1270) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:697) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:697) 在 scala.Option.foreach(Option.scala:236) 在 org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:697) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:1496) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1458) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1447) 在 org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:48) 在 org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:567) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1824) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1837) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1850) 在 org.apache.spark.sql.execution.SparkPlan.executeTake(SparkPlan.scala:215) 在 org.apache.spark.sql.execution.Limit.executeCollect(basicOperators.scala:207) 在 org.apache.spark.sql.DataFrame$$anonfun$collect$1.apply(DataFrame.scala:1385) 在 org.apache.spark.sql.DataFrame$$anonfun$collect$1.apply(DataFrame.scala:1385) 在 org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:56) 在 org.apache.spark.sql.DataFrame.withNewExecutionId(DataFrame.scala:1903) 在 org.apache.spark.sql.DataFrame.collect(DataFrame.scala:1384) 在 org.apache.spark.sql.DataFrame.head(DataFrame.scala:1314) 在 org.apache.spark.sql.DataFrame.take(DataFrame.scala:1377) 在 org.apache.spark.sql.DataFrame.showString(DataFrame.scala:178) 在 org.apache.spark.sql.DataFrame.show(DataFrame.scala:401) 在 org.apache.spark.sql.DataFrame.show(DataFrame.scala:362) 在 org.apache.spark.sql.DataFrame.show(DataFrame.scala:370) 在 comScore.DayWiseDeviceIDMatch$.main(DayWiseDeviceIDMatch.scala:62) 在 comScore.DayWiseDeviceIDMatch.main(DayWiseDeviceIDMatch.scala) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57) 在 sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.lang.reflect.Method.invoke(Method.java:606) 在 org.apache.spark.deploy.SparkSubmit$.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:674) 在 org.apache.spark.deploy.SparkSubmit$.doRunMain$1(SparkSubmit.scala:180) 在 org.apache.spark.deploy.SparkSubmit$.submit(SparkSubmit.scala:205) 在 org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:120) 在 org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
【问题讨论】:
-
您在哪里运行此作业?本地还是集群?
-
在亚马逊 EMR 集群中,它有 200GB 内存
标签: scala apache-spark apache-spark-sql