【发布时间】:2016-04-24 12:17:25
【问题描述】:
我已将 ELK 与 Pyspark 集成。
将RDD作为ELK数据保存在本地文件系统上
rdd.saveAsTextFile("/tmp/ELKdata")
logData = sc.textFile('/tmp/ELKdata/*')
errors = logData.filter(lambda line: "raw1-VirtualBox" in line)
errors.count()
我得到的值是 35
errors.first()
我得到了输出
(u'AVI0UK0KZsowGuTwoQnN', {u'host': u'raw1-VirtualBox', u'ident': u'NetworkManager', u'pid': u'748', u'message': u" ( eth0): 设备状态改变: ip-config -> 辅助设备 (原因 'none') [70 90 0]", u'@timestamp': u'2016-01-12T10:59:48+05:30'})
当我尝试在 pyspark 的弹性搜索中写入数据时,出现错误
errors.saveAsNewAPIHadoopFile(
path='-',
outputFormatClass="org.elasticsearch.hadoop.mr.EsOutputFormat",
keyClass="org.apache.hadoop.io.NullWritable",
valueClass="org.elasticsearch.hadoop.mr.LinkedMapWritable",
conf= {"es.resource" : "logstash-2016.01.12/errors})
巨大的java错误
org.apache.spark.SparkException:不能使用 java.lang.String 类型的 RDD 元素 在 org.apache.spark.api.python.SerDeUtil$$anonfun$pythonToPairRDD$1$$anonfun$apply$3.apply(SerDeUtil.scala:113) 在 org.apache.spark.api.python.SerDeUtil$$anonfun$pythonToPairRDD$1$$anonfun$apply$3.apply(SerDeUtil.scala:108) 在 scala.collection.Iterator$$anon$11.next(Iterator.scala:328) 在 scala.collection.Iterator$$anon$11.next(Iterator.scala:328) 在 org.apache.spark.rdd.PairRDDFunctions$$anonfun$12.apply(PairRDDFunctions.scala:921) 在 org.apache.spark.rdd.PairRDDFunctions$$anonfun$12.apply(PairRDDFunctions.scala:903) 在 org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:62) 在 org.apache.spark.scheduler.Task.run(Task.scala:54) 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:177) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615) 在 java.lang.Thread.run(Thread.java:745) 16/01/12 17:20:13 INFO TaskSetManager:在阶段 31.0 启动任务 1.0(TID 62,localhost,PROCESS_LOCAL,1181 字节) 16/01/12 17:20:13 WARN TaskSetManager:在阶段 31.0(TID 61,本地主机)中丢失任务 0.0:org.apache.spark.SparkException:无法使用 java.lang.String 类型的 RDD 元素 org.apache.spark.api.python.SerDeUtil$$anonfun$pythonToPairRDD$1$$anonfun$apply$3.apply(SerDeUtil.scala:113) org.apache.spark.api.python.SerDeUtil$$anonfun$pythonToPairRDD$1$$anonfun$apply$3.apply(SerDeUtil.scala:108) scala.collection.Iterator$$anon$11.next(Iterator.scala:328) scala.collection.Iterator$$anon$11.next(Iterator.scala:328) org.apache.spark.rdd.PairRDDFunctions$$anonfun$12.apply(PairRDDFunctions.scala:921) org.apache.spark.rdd.PairRDDFunctions$$anonfun$12.apply(PairRDDFunctions.scala:903) org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:62) org.apache.spark.scheduler.Task.run(Task.scala:54) org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:177) java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145) java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615) java.lang.Thread.run(Thread.java:745) 16/01/12 17:20:13 ERROR TaskSetManager: 阶段 31.0 中的任务 0 失败 1 次;中止工作 16/01/12 17:20:13 INFO TaskSchedulerImpl: 取消阶段 31 16/01/12 17:20:13 INFO TaskSchedulerImpl:第 31 阶段已取消 2012 年 16 月 1 日 17:20:13 信息执行器:执行器正在尝试在阶段 31.0 (TID 62) 中终止任务 1.0 2012 年 16 月 1 日 17:20:13 信息 DAGScheduler:无法在 PythonRDD.scala:665 处运行 saveAsNewAPIHadoopFile 回溯(最近一次通话最后): 文件“”,第 6 行,在 文件“/opt/spark/python/pyspark/rdd.py”,第 1213 行,在 saveAsNewAPIHadoopFile 中 keyConverter, valueConverter, jconf) 文件“/opt/spark/python/lib/py4j-0.8.2.1-src.zip/py4j/java_gateway.py”,第 538 行,在 __call__ 文件“/opt/spark/python/lib/py4j-0.8.2.1-src.zip/py4j/protocol.py”,第 300 行,在 get_return_value py4j.protocol.Py4JJavaError16/01/12 17:20:13 INFO 执行器:在阶段 31.0 (TID 62) 中运行任务 1.0 2012 年 16 月 1 日 17:20:13 错误执行程序:阶段 31.0 中的任务 1.0 异常(TID 62) org.apache.spark.TaskKilledException 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:168) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615) 在 java.lang.Thread.run(Thread.java:745) 16/01/12 17:20:13 WARN TaskSetManager:在 31.0 阶段丢失任务 1.0(TID 62,本地主机):org.apache.spark.TaskKilledException: org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:168) java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145) java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615) java.lang.Thread.run(Thread.java:745) 2012 年 16 月 1 日 17:20:13 信息 TaskSchedulerImpl:从池中删除了任务已全部完成的 TaskSet 31.0 : 调用 z:org.apache.spark.api.python.PythonRDD.saveAsNewAPIHadoopFile 时出错。 :org.apache.spark.SparkException:作业因阶段失败而中止:阶段 31.0 中的任务 0 失败 1 次,最近一次失败:阶段 31.0 中丢失任务 0.0(TID 61,本地主机):org.apache.spark.SparkException:不能使用 java.lang.String 类型的 RDD 元素 org.apache.spark.api.python.SerDeUtil$$anonfun$pythonToPairRDD$1$$anonfun$apply$3.apply(SerDeUtil.scala:113) org.apache.spark.api.python.SerDeUtil$$anonfun$pythonToPairRDD$1$$anonfun$apply$3.apply(SerDeUtil.scala:108) scala.collection.Iterator$$anon$11.next(Iterator.scala:328) scala.collection.Iterator$$anon$11.next(Iterator.scala:328) org.apache.spark.rdd.PairRDDFunctions$$anonfun$12.apply(PairRDDFunctions.scala:921) org.apache.spark.rdd.PairRDDFunctions$$anonfun$12.apply(PairRDDFunctions.scala:903) org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:62) org.apache.spark.scheduler.Task.run(Task.scala:54) org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:177) java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145) java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615) java.lang.Thread.run(Thread.java:745) 驱动程序堆栈跟踪: 在 org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1185) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1174) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1173) 在 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:1173) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:688) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:688) 在 scala.Option.foreach(Option.scala:236) 在 org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:688) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessActor$$anonfun$receive$2.applyOrElse(DAGScheduler.scala:1391) 在 akka.actor.ActorCell.receiveMessage(ActorCell.scala:498) 在 akka.actor.ActorCell.invoke(ActorCell.scala:456) 在 akka.dispatch.Mailbox.processMailbox(Mailbox.scala:237) 在 akka.dispatch.Mailbox.run(Mailbox.scala:219) 在 akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(AbstractDispatcher.scala:386) 在 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)如果我手动完成,则可以写入数据
errors = logData.filter(lambda line: "raw1-VirtualBox" in line)
errors = errors.map(lambda item: ('AVI0UK0KZsowGuTwoQnP',{"host": "raw1-VirtualBox",
"ident": "NetworkManager",
"pid": "69",
"message": " sucess <info> (eth0): device state change: ip-config -> secondaries (reason 'none') [70 90 0]",
"@timestamp": "2016-01-12T10:59:48+05:30"
}))
但我想在弹性搜索中写入过滤数据和托管数据。
【问题讨论】:
标签: elasticsearch apache-spark pyspark elastic-map-reduce