首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >Spark Streaming kafka concurrentModificationException

Spark Streaming kafka concurrentModificationException
EN

Stack Overflow用户
提问于 2017-12-03 22:01:38
回答 0查看 1.2K关注 0票数 2

我使用的是Spark流媒体应用程序。应用程序使用直接流从Kafka topic (具有200个分区)中读取消息。应用程序偶尔会抛出ConcurrentModificationException->

代码语言:javascript
复制
java.util.ConcurrentModificationException: KafkaConsumer is not safe for multi-threaded access
at org.apache.kafka.clients.consumer.KafkaConsumer.acquire(KafkaConsumer.java:1431)
at org.apache.kafka.clients.consumer.KafkaConsumer.close(KafkaConsumer.java:1361)
at org.apache.spark.streaming.kafka010.CachedKafkaConsumer$$anon$1.removeEldestEntry(CachedKafkaConsumer.scala:128)
at java.util.LinkedHashMap.afterNodeInsertion(LinkedHashMap.java:299)
at java.util.HashMap.putVal(HashMap.java:663)
at java.util.HashMap.put(HashMap.java:611)
at org.apache.spark.streaming.kafka010.CachedKafkaConsumer$.get(CachedKafkaConsumer.scala:158)
at org.apache.spark.streaming.kafka010.KafkaRDD$KafkaRDDIterator.<init>(KafkaRDD.scala:211)
at org.apache.spark.streaming.kafka010.KafkaRDD.compute(KafkaRDD.scala:186)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:323)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:287)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
at org.apache.spark.scheduler.Task.run(Task.scala:99)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:282)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
at java.lang.Thread.run(Thread.java:745)

我的spark集群有两个节点。Spark版本是2.1。应用程序运行两个执行器。从我从异常和kafka消费者代码中可以看出,似乎同一个kakfa消费者正在被两个线程使用。我不知道为什么两个线程访问同一个接收器。理想情况下,每个执行器都应该有一个独占的kafka接收器,由单个线程提供服务,该线程必须读取所有分配的分区的消息。从kafka读取的代码片段->

代码语言:javascript
复制
JavaInputDStream<ConsumerRecord<String, String>> consumerRecords = KafkaUtils.createDirectStream(
jssc,
LocationStrategies.PreferConsistent(),
ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams));
EN

回答

页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/47619148

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档