【问题标题】:Kafka Spark Streaming Error - java.lang.NoClassDefFoundError: org/apache/spark/sql/connector/read/streaming/ReportsSourceMetricsKafka Spark 流式处理错误 - java.lang.NoClassDefFoundError: org/apache/spark/sql/connector/read/streaming/ReportsSourceMetrics
【发布时间】:2022-01-27 02:55:12
【问题描述】:

我正在使用 Spark 3.1.2、Kafka 2.8.1 和 Scala 2.12.1

在集成 Kafka 和 Spark 流时遇到错误 -

java.lang.NoClassDefFoundError: org/apache/spark/sql/connector/read/streaming/ReportsSourceMetrics

具有依赖关系的 Spark-shell 命令 - spark-shell --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2

 org.apache.spark#spark-sql-kafka-0-10_2.12 added as a dependency
 :: resolving dependencies :: org.apache.spark#spark-submit-parent-3643b83d-a2f8-43d1-941f-a125272f3905;1.0
         confs: [default]
         found org.apache.spark#spark-sql-kafka-0-10_2.12;3.1.2 in central
         found org.apache.spark#spark-token-provider-kafka-0-10_2.12;3.1.2 in central
         found org.apache.kafka#kafka-clients;2.6.0 in central
         found com.github.luben#zstd-jni;1.4.8-1 in central
         found org.lz4#lz4-java;1.7.1 in central
         found org.xerial.snappy#snappy-java;1.1.8.2 in central
         found org.slf4j#slf4j-api;1.7.30 in central
         found org.spark-project.spark#unused;1.0.0 in central
         found org.apache.commons#commons-pool2;2.6.2 in central
 :: resolution report :: resolve 564ms :: artifacts dl 9ms
         :: modules in use:
         com.github.luben#zstd-jni;1.4.8-1 from central in [default]
         org.apache.commons#commons-pool2;2.6.2 from central in [default]
         org.apache.kafka#kafka-clients;2.6.0 from central in [default]
         org.apache.spark#spark-sql-kafka-0-10_2.12;3.1.2 from central in [default]
         org.apache.spark#spark-token-provider-kafka-0-10_2.12;3.1.2 from central in [default]
         org.lz4#lz4-java;1.7.1 from central in [default]
         org.slf4j#slf4j-api;1.7.30 from central in [default]
         org.spark-project.spark#unused;1.0.0 from central in [default]
         org.xerial.snappy#snappy-java;1.1.8.2 from central in [default]
         ---------------------------------------------------------------------
         |                  |            modules            ||   artifacts   |
         |       conf       | number| search|dwnlded|evicted|| number|dwnlded|
         ---------------------------------------------------------------------
         |      default     |   9   |   0   |   0   |   0   ||   9   |   0   |
         ---------------------------------------------------------------------
 :: retrieving :: org.apache.spark#spark-submit-parent-3643b83d-a2f8-43d1-941f-a125272f3905
         confs: [default]
         0 artifacts copied, 9 already retrieved (0kB/15ms)
 21/12/28 17:46:21 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
 Using Spark's default log4j profile: org/apache/spark/log4j-defaults.properties
 Setting default log level to "WARN".
 To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
 21/12/28 17:46:28 WARN Utils: Service 'SparkUI' could not bind on port 4040. Attempting port 4041.
 Spark context Web UI available at http://*******:4041
 Spark context available as 'sc' (master = local[*], app id = local-1640693788919).
 Spark session available as 'spark'.
 Welcome to
       ____              __
      / __/__  ___ _____/ /__
     _\ \/ _ \/ _ `/ __/  '_/
    /___/ .__/\_,_/_/ /_/\_\   version 3.1.2
       /_/
 
 Using Scala version 2.12.10 (OpenJDK 64-Bit Server VM, Java 1.8.0_292)
 Type in expressions to have them evaluated.
 Type :help for more information.
    
    val df = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "127.0.1.1:9092").option("subscribe", "Topic").option("startingOffsets", "earliest").load()
    
    df.printSchema()
    
    import org.apache.spark.sql.types._
    val schema = new StructType().add("id",IntegerType).add("fname",StringType).add("lname",StringType)
    val personStringDF = df.selectExpr("CAST(value AS STRING)")
    val personDF = personStringDF.select(from_json(col("value"), schema).as("data")).select("data.*")
     
    personDF.writeStream.format("console").outputMode("append").start().awaitTermination()
    
    Exception in thread "stream execution thread for [id = 44e8f8bf-7d94-4313-9d2b-88df8f5bc10f, runId = 3b4c63c4-9062-4288-a681-7dd6cfb836d0]" java.lang.NoClassDefFoundError: org/apache/spark/sql/connector/read/streaming/ReportsSourceMetrics

【问题讨论】:

    标签: scala apache-spark apache-spark-sql spark-structured-streaming


    【解决方案1】:

    我遇到了几乎相同的问题 - 相同的异常,但在 spark-submit 中。 我通过将 Spark 升级到版本3.2.0 解决了这个问题。 我还使用了org.apache.spark:spark-sql-kafka-0-10_2.123.2.0 版本,完整的命令是:

    $SPARK_HOME/bin/spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0 script.py
    

    【讨论】:

    • 感谢您的评论。我会升级到最新的试试看。
    【解决方案2】:

    Spark_version 3.1.2

    Scala_version 2.12.10

    Kafka_version 2.8.1

    注意:当我们将--packages org.apache.spark:spark-sql-kafka-0-10_2.12:V.V.Vspark-shellspark-submit 一起使用时,版本非常重要。在哪里(V.V.V = Spark_version)

    我按照spark-kafka-example 给出的以下步骤进行操作:

    1. 启动生产者 $ kafka-console-producer.sh --broker-list Kafka-Server-IP:9092 --topic kafka-spark-test

    您应该会在控制台上看到提示符>。输入一些生产者的测试数据。

    >{"name":"foo","dob_year":1995,"gender":"M","salary":2000}
    >{"name":"bar","dob_year":1996,"gender":"M","salary":2500}
    >{"name":"baz","dob_year":1997,"gender":"F","salary":3500}
    >{"name":"foo-bar","dob_year":1998,"gender":"M","salary":4000}
    
    
    1. 按如下方式启动spark-shell
    spark-shell --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2
    

    注意:我使用的是 3.1.2。成功启动后,您将看到类似以下内容:

    Spark session available as 'spark'.
    Welcome to
          ____              __
         / __/__  ___ _____/ /__
        _\ \/ _ \/ _ `/ __/  '_/
       /___/ .__/\_,_/_/ /_/\_\   version 3.1.2
          /_/
             
    Using Scala version 2.12.10 (OpenJDK 64-Bit Server VM, Java 11.0.13)
    Type in expressions to have them evaluated.
    Type :help for more information.
    
    
    1. 输入导入并创建 DataFrame,然后打印架构。
    val df = spark.readStream.
          format("kafka"). 
          option("kafka.bootstrap.servers", "Kafka-Server-IP:9092").
          option("subscribe", "kafka-spark-test").
          option("startingOffsets", "earliest").
          load()
    
    df.printSchema()
    
    1. 成功执行应产生以下结果:
    scala> df.printSchema()
    root
     |-- key: binary (nullable = true)
     |-- value: binary (nullable = true)
     |-- topic: string (nullable = true)
     |-- partition: integer (nullable = true)
     |-- offset: long (nullable = true)
     |-- timestamp: timestamp (nullable = true)
     |-- timestampType: integer (nullable = true)
    
    
    
    1. 将 DataFrame 的二进制值转换为字符串。我正在显示带有输出的命令
    scala>      val personStringDF = df.selectExpr("CAST(value AS STRING)")
    
    
    personStringDF: org.apache.spark.sql.DataFrame = [value: string]
    
    
    1. DataFrame 的品牌和架构。我正在显示带有输出的命令
    scala> val schema = new StructType().
         |       add("name",StringType).
         |       add("dob_year",IntegerType).
         |       add("gender",StringType).
         |       add("salary",IntegerType)
    
    
    schema: org.apache.spark.sql.types.StructType = StructType(StructField(name,StringType,true), StructField(dob_year,IntegerType,true), StructField(gender,StringType,true), StructField(salary,IntegerType,true))
    
    
    1. 选择数据
    scala>  val personDF = personStringDF.select(from_json(col("value"), schema).as("data")).select("data.*")
    
    personDF: org.apache.spark.sql.DataFrame = [name: string, dob_year: int ... 2 more fields]
    
    
    1. 在控制台上写入流
    
    scala>  personDF.writeStream.
         |       format("console").
         |       outputMode("append").
         |       start().
         |       awaitTermination()
    
    
    

    您将看到以下输出:

    -------------------------------------------                                     
    Batch: 0
    -------------------------------------------
    +-------+--------+------+------+
    |   name|dob_year|gender|salary|
    +-------+--------+------+------+
    |    foo|    1981|     M|  2000|
    |    bar|    1982|     M|  2500|
    |    baz|    1983|     F|  3500|
    |foo-bar|    1984|     M|  4000|
    +-------+--------+------+------+
    
    

    如果你的 kafka producer 还在运行,你可以输入一个新行,每次在 producer 中输入新数据时,你都会在 Batch: 1 中看到新数据。

    这是我们从控制台生产者输入数据并在火花控制台中消费的典型示例。

    祝你好运! :)

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-04-06
      • 2018-03-10
      • 1970-01-01
      • 2020-10-14
      • 1970-01-01
      • 2019-06-07
      相关资源
      最近更新 更多