【发布时间】:2019-04-11 15:56:38
【问题描述】:
我无法使用 spark-streaming 运行 Kafka。以下是我到目前为止所采取的步骤:
下载
jar文件“spark-streaming-kafka-0-8-assembly_2.10-2.2.0.jar”并移至/home/ec2-user/spark-2.0.0-bin-hadoop2.7/jars将此行添加到
/home/ec2-user/spark-2.0.0-bin-hadoop2.7/conf/spark-defaults.conf.template->spark.jars.packages org.apache.spark:spark-streaming-kafka-0-8-assembly_2.10:2.2.0
卡夫卡版本:kafka_2.10-0.10.2.2
Jar 文件版本:spark-streaming-kafka-0-8-assembly_2.10-2.2.0.jar
Python 代码:
os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-streaming-kafka-0-8-assembly_2.10-2.2.0 pyspark-shell'
kvs = KafkaUtils.createDirectStream(ssc, ["divolte-data"], {"metadata.broker.list": "localhost:9092"})
但我仍然收到以下错误:
Py4JJavaError: An error occurred while calling o39.createDirectStreamWithoutMessageHandler.
: java.lang.NoClassDefFoundError: Could not initialize class kafka.consumer.FetchRequestAndResponseStatsRegistry$
at kafka.consumer.SimpleConsumer.<init>(SimpleConsumer.scala:39)
at org.apache.spark.streaming.kafka.KafkaCluster.connect(KafkaCluster.scala:59)
我做错了什么?
【问题讨论】:
-
你是如何配置pom的?你在使用 `metrics-core-2.2.0.jar
?spark-shell --.jars metrics-core-2.2.0.jar` -
您使用的是
spark-2.0.0,但您的罐子是为2.2.0... 那些版本应该是一样的
标签: java apache-spark pyspark apache-kafka spark-streaming