【问题标题】:using spark streaming from pyspark 2.4.4使用来自 pyspark 2.4.4 的火花流
【发布时间】:2020-04-18 21:12:41
【问题描述】:

我在 k8s 容器中设置了 spark 2.4.4 版本。我正在尝试为使用这样的火花流编写一个简单的 hello world:

from pyspark import SparkContext
from pyspark.sql import SparkSession
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils
spark = SparkSession.builder.appName("pyspark-kafka").getOrCreate()
sc.setLogLevel("WARN")

ssc = StreamingContext(sc, 60)
kafkaStream = KafkaUtils.createDirectStream(ssc, ['users-update'], {"metadata.broker.list":'pubsub-0.pubsub:9092,pubsub-1.pubsub:9092,pubsub-2.pubsub:9092'})

请注意,pubsub-x.pubsub 是我的容器可见的 kafka 代理。 (在我的最后一行 pyspark 代码中,一个简单的 python 程序直接使用带有代理和主题的 kafka-python 客户端工作得很好。)

我收到此错误消息:

________________________________________________________________________________________________

  Spark Streaming's Kafka libraries not found in class path. Try one of the following.

  1. Include the Kafka library and its dependencies with in the
     spark-submit command as

     $ bin/spark-submit --packages org.apache.spark:spark-streaming-kafka-0-8:2.4.4 ...

  2. Download the JAR of the artifact from Maven Central http://search.maven.org/,
     Group Id = org.apache.spark, Artifact Id = spark-streaming-kafka-0-8-assembly, Version = 2.4.4.
     Then, include the jar in the spark-submit command as

     $ bin/spark-submit --jars <spark-streaming-kafka-0-8-assembly.jar> ...

________________________________________________________________________________________________

maven 上的任何地方都没有 2.4.4 版本的 kafka 库。 https://search.maven.org/search?q=spark%20kafka 显示最后发布的 jar 是 2.10 或 2.11 版本。

我的 pyspark 安装中确实有一个 spark-streaming_2.12-2.4.4.jar jar,但它似乎没有正确的 kafka 类。

感谢您的任何指点! --斯里达尔

【问题讨论】:

    标签: apache-spark pyspark apache-kafka


    【解决方案1】:

    Spark v2.4.4 是使用 scala v2.11 预构建的。从火花下载页面:

    请注意,Spark 是使用 Scala 2.11 预构建的,但版本 2.4.2 是使用 Scala 2.12 预构建的。

    所以,基本上2.102.11 是构建 spark 的 scala 版本,你应该下载 spark-streaming-kafka jar,它是在你的情况下使用相同版本的 scala 构建的 2.11

    我已经检查了 spark 2.4.4 中的 jars 文件夹,并且那里存在 spark-streaming_2.11-2.4.4.jar jar。因此,如果您已将 spark-streaming_2.12-2.4.4.jar 从外部添加到类路径中,则应删除它,否则您将得到版本不匹配。

    您可以从here 下载spark-streaming-kafka-0-8-assembly.jar 而且我认为您还需要从here 添加kafka-clients jar。

    【讨论】:

    • 2.4.4 也是使用 Scala 2.12 构建的 - spark-2.4.4-bin-without-hadoop-scala-2.12.tgz
    • 感谢您的建议。我通过运行pip install kafka-python 安装了kafka python 客户端之后,我的小python sn-p 连接到kafka 工作正常。我根据您的建议安装了spark-streaming-kafka-0-8-assembly_2.11-2.4.4.jarspark-streaming_2.11-2.4.4.jar,并确保jar 中没有其他流媒体jar。我仍然收到错误,An error occurred while calling o29.load. :org.apache.spark.sql.AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming...
    • 我更新的产生错误的pyspark代码sn-p是:from pyspark import SparkContext from pyspark.sql import SparkSession spark = SparkSession.builder.appName("pyspark-kafka").getOrCreate() sc.setLogLevel("WARN") df = spark.readStream.format("kafka") \ .option("kafka.bootstrap.servers", "pubsub-0.pubsub:9092,pubsub-1.pubsub:9092,pubsub-2.pubsub:9092") \ .option("subscribe", "users-update") \ .load()
    • 你能告诉我你的 spark-submit 命令吗?您是否将 kafka-clients jar 添加到我在上一行提到的类路径中?
    • 获胜配置为:spark-streaming-kafka-0-10-assembly_2.11-2.4.4.jarspark-sql-kafka-0-10_2.12-2.4.4.jar 以及 spark-streaming_2.11-2.4.4.jarkafka-clients-2.4.0.jar。感谢您的所有帮助!
    【解决方案2】:

    我的 pyspark 安装中确实有一个 spark-streaming_2.12-2.4.4.jar jar,但它似乎没有正确的 kafka 类。

    这只是 Spark 的基本 Streaming 包。 Spark 不附带 Kafka 类

    Spark Streaming 已被弃用,取而代之的是 Spark Structured Streaming

    你想要这个包用于带有 Scala 2.12 的 Spark

    'org.apache.spark:spark-sql-kafka-0-10_2.12:2.4.4'
    

    你会像这样开始,包括引导服务器的选项

    df = spark.readStream().format("kafka")
    

    【讨论】:

    • 您已经提供了 kafka 0.10 的依赖项,我认为根据错误 kafka 0.8 依赖项是必需的。
    • @wypul 错误只是存在,因为它是 from pyspark.streaming.kafka import KafkaUtils 没有可用软件包的一部分。 Structured Streaming 包与旧的流 API 具有 API 兼容性
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-08-21
    • 1970-01-01
    • 2018-08-09
    • 2020-10-25
    • 2016-06-22
    相关资源
    最近更新 更多