【发布时间】:2018-10-23 10:46:18
【问题描述】:
尝试用 spqrk 编写简单的程序。我必须按 myStruct 的一个属性对我的数据进行分组 - LogData:
public class LogData {
public String m_Host;
public String m_Timestamp;
public String m_Request;
public Integer m_Reply;
public String m_ColumnByteReply;
}
我尝试了什么:
JavaPairRDD <String, Iterable<LogData>> tmp = parsedData.groupBy(logData -> logData.m_Host);
JavaPairRDD<String, Iterable<LogData>> groupMap = parsedData.groupBy(new Function<LogData, String>() {
@Override
public String call(LogData logData) throws Exception {
return logData.m_Host;
}
});
而且很简单:
JavaPairRDD <String, Iterable<LogData>> tmp = parsedData.groupBy(logData -> logData.m_Host);
当我尝试输出 resultData 时,我的程序失败了。
错误:
18/05/13 20:51:32 错误执行程序:阶段 1.0 中任务 0.0 中的异常 (第 1 次) java.io.NotSerializableException:日志数据 在 java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1184) 在 java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548) 在 java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509) 在 java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432) 在 java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178) 在 java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348) 在 org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:42) 在 org.apache.spark.storage.DiskBlockObjectWriter.write(BlockObjectWriter.scala:195) 在 org.apache.spark.util.collection.ExternalSorter.spillToPartitionFiles(ExternalSorter.scala:370) 在 org.apache.spark.util.collection.ExternalSorter.insertAll(ExternalSorter.scala:211) 在 org.apache.spark.shuffle.sort.SortShuffleWriter.write(SortShuffleWriter.scala:65) 在 org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:68) 在 org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:41) 在 org.apache.spark.scheduler.Task.run(Task.scala:56) 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:196) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) 在 java.lang.Thread.run(Thread.java:748) 18/05/13 20:51:32 错误 TaskSetManager:阶段 1.0(TID 1)中的任务 0.0 具有不可序列化的结果:LogData;不重试 2013 年 5 月 18 日 20:51:32 信息 TaskSchedulerImpl:从池中删除了任务已全部完成的 TaskSet 1.0 18/05/13 20:51:32 INFO TaskSchedulerImpl:取消阶段 1 2013 年 5 月 18 日 20:51:32 信息 DAGScheduler:作业 1 失败:WordCount.java:92 计数,耗时 0,088518 秒 org.apache.spark.SparkException:作业因阶段失败而中止:阶段 1.0(TID 1)中的任务 0.0 具有不可序列化的结果:LogData 在 org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1214) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1203) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1202) 在 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:1202) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:696) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:696) 在 scala.Option.foreach(Option.scala:236) 在 org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:696) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessActor$$anonfun$receive$2.applyOrElse(DAGScheduler.scala:1420) 在 akka.actor.Actor$class.aroundReceive(Actor.scala:465) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessActor.aroundReceive(DAGScheduler.scala:1375) 在 akka.actor.ActorCell.receiveMessage(ActorCell.scala:516) 在 akka.actor.ActorCell.invoke(ActorCell.scala:487) 在 akka.dispatch.Mailbox.processMailbox(Mailbox.scala:238) 在 akka.dispatch.Mailbox.run(Mailbox.scala:220) 在 akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(AbstractDispatcher.scala:393) 在 scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) 在 scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339) 在 scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) 在 scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)
帮助:) 感谢您的回答!
【问题讨论】:
标签: java apache-spark hadoop