【问题标题】:Kafka with spark streaming integration error卡夫卡火花流集成错误
【发布时间】:2019-04-11 15:56:38
【问题描述】:

我无法使用 spark-streaming 运行 Kafka。以下是我到目前为止所采取的步骤:

  1. 下载jar文件“spark-streaming-kafka-0-8-assembly_2.10-2.2.0.jar”并移至/home/ec2-user/spark-2.0.0-bin-hadoop2.7/jars

  2. 将此行添加到/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


【解决方案1】:

spark-defaults.conf.template 只是一个模板,不会被 Spark 读取,因此不会加载您的 JAR。您必须复制/重命名此文件以删除模板后缀

如果您想使用这些特定的 JAR 文件,您还需要下载 Spark 2.2。

如果您要使用 Kafka 包,请确保您的 Spark 版本使用 Scala 2.10。否则使用2.11版本

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-08-15
    • 2016-08-03
    • 2018-09-15
    • 2018-08-13
    • 2018-02-24
    • 2023-03-19
    • 2016-12-19
    • 1970-01-01
    相关资源
    最近更新 更多