我有一个
Scala
Spark Streaming
从同一主题接收来自3个不同主题的数据的应用程序
Kafka producers
.
Spark streaming应用程序位于带有主机的计算机上
0.0.0.179
,Kafka服务器位于具有主机的计算机上
0.0.0.178
这个
卡夫卡制作人
在机器上,
0.0.0.180
,
0.0.0.181
,
0.0.0.182
.
当我尝试运行
火花流
应用程序出现以下错误
线程“main”组织中出现异常。阿帕奇。火花SparkException:作业
由于阶段失败而中止:阶段19.0中的任务0失败1次,
最近的失败:在阶段19.0中丢失了任务0.0(TID 19,localhost):
Java语言util。ConcurrentModificationException:卡夫卡消费者不安全
用于多线程访问
在
组织。阿帕奇。卡夫卡。客户。消费者卡夫卡康萨默尔。seek(KafkaConsumer.java:1198)
在
组织。阿帕奇。火花流动。kafka010.CachedKafkaConsumer。seek(CachedKafkaConsumer。scala:95)
在
组织。阿帕奇。火花流动。kafka010.CachedKafkaConsumer。获取(CachedKafkaConsumer.scala:69)
在
组织。阿帕奇。火花流动。kafka010.kafkard$kafkarditerator。下一个(卡夫卡德·斯卡拉:228)
在
组织。阿帕奇。火花流动。kafka010.kafkard$kafkarditerator。下一个(卡夫卡德·斯卡拉:194)
在scala。收集迭代器$$不超过11美元。下一步(迭代器.scala:409)位于
斯卡拉。收集迭代器$$不超过11美元。下一步(迭代器.scala:409)位于
组织。阿帕奇。火花rdd。pairddfunctions$$anonfun$saveAsHadoopDataset$1$$anonfun$13$$anonfun$apply$7。应用$mcV$sp(pairrdFunctions.scala:1204)
在
组织。阿帕奇。火花rdd。pairddfunctions$$anonfun$saveAsHadoopDataset$1$$anonfun$13$$anonfun$apply$7。应用(pairddfunctions.scala:1203)
在
组织。阿帕奇。火花rdd。pairddfunctions$$anonfun$saveAsHadoopDataset$1$$anonfun$13$$anonfun$apply$7。应用(pairddfunctions.scala:1203)
在
组织。阿帕奇。火花util。Utils美元。尝试使用SafeFinallyandFailureCallbacks(实用规模:1325)
在
组织。阿帕奇。火花rdd。PairRDDFunctions$$anonfun$SaveAshadopDataSet$1$$anonfun$13。应用(pairddfunctions.scala:1211)
在
组织。阿帕奇。火花rdd。PairRDDFunctions$$anonfun$SaveAshadopDataSet$1$$anonfun$13。应用(pairddfunctions.scala:1190)
位于组织。阿帕奇。火花调度程序。结果任务。runTask(ResultTask.scala:70)
位于组织。阿帕奇。火花调度程序。任务运行(任务.scala:85)
组织。阿帕奇。火花执行人。执行者$TaskRunner。运行(执行器scala:274)
在
Java语言util。同时发生的线程池执行器。runWorker(ThreadPoolExecutor.java:1142)
在
Java语言util。同时发生的ThreadPoolExecutor$工作者。运行(ThreadPoolExecutor.java:617)
在java。lang.Thread。运行(Thread.java:748)
现在我读了数千篇不同的帖子,但似乎没有人能找到解决这个问题的方法。
如何在我的应用程序中处理此问题?我是否必须修改Kakfa的一些参数(目前
num.partition
参数设置为1)?
以下是我的申请代码:
// Create the context with a 5 second batch size
val sparkConf = new SparkConf().setAppName("SparkScript").set("spark.driver.allowMultipleContexts", "true").set("spark.streaming.concurrentJobs", "3").setMaster("local[4]")
val sc = new SparkContext(sparkConf)
val ssc = new StreamingContext(sc, Seconds(3))
case class Thema(name: String, metadata: String)
case class Tempo(unit: String, count: Int, metadata: String)
case class Spatio(unit: String, metadata: String)
case class Stt(spatial: Spatio, temporal: Tempo, thematic: Thema)
case class Location(latitude: Double, longitude: Double, name: String)
case class Datas1(location : Location, timestamp : String, windspeed : Double, direction: String, strenght : String)
case class Sensors1(sensor_name: String, start_date: String, end_date: String, data1: Datas1, stt: Stt)
val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "0.0.0.178:9092",
"key.deserializer" -> classOf[StringDeserializer].getCanonicalName,
"value.deserializer" -> classOf[StringDeserializer].getCanonicalName,
"group.id" -> "test_luca",
"auto.offset.reset" -> "earliest",
"enable.auto.commit" -> (false: java.lang.Boolean)
)
val topics1 = Array("topics1")
val s1 = KafkaUtils.createDirectStream[String, String](ssc, PreferConsistent, Subscribe[String, String](topics1, kafkaParams)).map(record => {
implicit val formats = DefaultFormats
parse(record.value).extract[Sensors1]
}
)
s1.print()
s1.saveAsTextFiles("results/", "")
ssc.start()
ssc.awaitTermination()
非常感谢。