代码之家  ›  专栏  ›  技术社区  ›  lu_ferra

Java语言util。ConcurrentModificationException:KafkaConsumer对于多线程访问不安全

  •  2
  • lu_ferra  · 技术社区  · 8 年前

    我有一个 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()
    

    非常感谢。

    2 回复  |  直到 8 年前
        1
  •  5
  •   Yuval Itzchakov    8 年前

    您的问题在于:

    s1.print()
    s1.saveAsTextFiles("results/", "")
    

    因为Spark创建了一个流图,您在这里定义了两个流:

    Read from Kafka -> Print to console
    Read from Kafka -> Save to text file
    

    Spark将尝试同时运行这两个图,因为它们彼此独立。由于Kafka使用缓存消费者方法,因此它实际上是在尝试对两个流执行使用相同的消费者。

    您可以做的是缓存 DStream 在运行这两个查询之前:

    val dataFromKafka = KafkaUtils.createDirectStream[String, String](ssc, PreferConsistent, Subscribe[String, String](topics1, kafkaParams)).map(/* stuff */)
    
    val cachedStream = dataFromKafka.cache()
    cachedStream.print()
    cachedStream.saveAsTextFiles("results/", "")
    
        2
  •  0
  •   kartik kudada Taha Naqvi    5 年前

    使用缓存对我很有用。在我的例子中,打印、转换然后在JavaPairDstream上打印给了我这个错误。 在第一次打印之前,我使用了缓存,它对我很有用。

    s1.print()
    s1.saveAsTextFiles("results/", "")
    

    下面的代码可以工作,我用过类似的代码。

    s1.cache();
    s1.print();
    s1.saveAsTextFiles("results/", "");
    
    推荐文章