【发布时间】:2016-03-21 04:44:58
【问题描述】:
我正在尝试设置 Spark Streaming 以从 Kafka 队列中获取消息。我收到以下错误:
py4j.protocol.Py4JJavaError: An error occurred while calling o30.createDirectStream.
: org.apache.spark.SparkException: java.nio.channels.ClosedChannelException
org.apache.spark.SparkException: Couldn't find leader offsets for Set([test-topic,0])
at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$checkErrors$1.apply(KafkaCluster.scala:366)
at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$checkErrors$1.apply(KafkaCluster.scala:366)
at scala.util.Either.fold(Either.scala:97)
这是我正在执行的代码 (pyspark):
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils
directKafkaStream = KafkaUtils.createDirectStream(ssc, ["test-topic"], {"metadata.broker.list": "host.domain:9092"})
ssc.start()
ssc.awaitTermination()
有几个类似的帖子有相同的错误。在所有情况下,原因都是空的 kafka 主题。我的“测试主题”中有消息。我可以把它们弄出来
kafka-console-consumer --zookeeper host.domain:2181 --topic test-topic --from-beginning --max-messages 100
有谁知道可能是什么问题?
我正在使用:
- Spark 1.5.2 (apache)
- Kafka 0.8.2.0+kafka1.3.0 (CDH 5.4.7)
【问题讨论】:
-
我认为这是缺少领导者的问题,请查看helpful i think
-
我也遇到了同样的问题,请问您找到解决方法了吗?我使用 spark 1.6.1 和 kafka 0.8.2.1
-
我将我的偏移量存储在 zookeeper 中。我清除/休息了我存储的偏移量,这个错误不再出现。
标签: apache-spark apache-kafka spark-streaming