【问题标题】:PySpark Kafka - NoClassDefFound: org/apache/commons/pool2PySpark Kafka - NoClassDefFound:org/apache/commons/pool2
【发布时间】:2021-09-14 06:11:27
【问题描述】:

我在将数据从 kafka 主题打印到控制台时遇到问题。 我收到的错误消息如下图所示。

如上图所示,在第 0 批之后,它不会进一步处理。

所有这些都是错误消息的快照。我不明白发生错误的根本原因。请帮帮我。

以下是kafka和spark版本:

spark version: spark-3.1.1-bin-hadoop2.7
kafka version: kafka_2.13-2.7.0

我正在使用以下罐子:

kafka-clients-2.7.0.jar 
spark-sql-kafka-0-10_2.12-3.1.1.jar 
spark-token-provider-kafka-0-10_2.12-3.1.1.jar 

这是我的代码:

spark = SparkSession \
        .builder \
        .appName("Pyspark structured streaming with kafka and cassandra") \
        .master("local[*]") \
        .config("spark.jars","file:///C://Users//shivani//Desktop//Spark//kafka-clients-2.7.0.jar,file:///C://Users//shivani//Desktop//Spark//spark-sql-kafka-0-10_2.12-3.1.1.jar,file:///C://Users//shivani//Desktop//Spark//spark-cassandra-connector-2.4.0-s_2.11.jar,file:///D://mysql-connector-java-5.1.46//mysql-connector-java-5.1.46.jar,file:///C://Users//shivani//Desktop//Spark//spark-token-provider-kafka-0-10_2.12-3.1.1.jar")\
        .config("spark.executor.extraClassPath","file:///C://Users//shivani//Desktop//Spark//kafka-clients-2.7.0.jar,file:///C://Users//shivani//Desktop//Spark//spark-sql-kafka-0-10_2.12-3.1.1.jar,file:///C://Users//shivani//Desktop//Spark//spark-cassandra-connector-2.4.0-s_2.11.jar,file:///D://mysql-connector-java-5.1.46//mysql-connector-java-5.1.46.jar,file:///C://Users//shivani//Desktop//Spark//spark-token-provider-kafka-0-10_2.12-3.1.1.jar")\
        .config("spark.executor.extraLibrary","file:///C://Users//shivani//Desktop//Spark//kafka-clients-2.7.0.jar,file:///C://Users//shivani//Desktop//Spark//spark-sql-kafka-0-10_2.12-3.1.1.jar,file:///C://Users//shivani//Desktop//Spark//spark-cassandra-connector-2.4.0-s_2.11.jar,file:///D://mysql-connector-java-5.1.46//mysql-connector-java-5.1.46.jar,file:///C://Users//shivani//Desktop//Spark//spark-token-provider-kafka-0-10_2.12-3.1.1.jar")\
        .config("spark.driver.extraClassPath","file:///C://Users//shivani//Desktop//Spark//kafka-clients-2.7.0.jar,file:///C://Users//shivani//Desktop//Spark//spark-sql-kafka-0-10_2.12-3.1.1.jar,file:///C://Users//shivani//Desktop//Spark//spark-cassandra-connector-2.4.0-s_2.11.jar,file:///D://mysql-connector-java-5.1.46//mysql-connector-java-5.1.46.jar,file:///C://Users//shivani//Desktop//Spark//spark-token-provider-kafka-0-10_2.12-3.1.1.jar")\
        .getOrCreate()
    spark.sparkContext.setLogLevel("ERROR")


#streaming dataframe that reads from kafka topic
    df_kafka=spark.readStream\
    .format("kafka")\
    .option("kafka.bootstrap.servers",kafka_bootstrap_servers)\
    .option("subscribe",kafka_topic_name)\
    .option("startingOffsets", "latest") \
    .load()

    print("Printing schema of df_kafka:")
    df_kafka.printSchema()

    #converting data from kafka broker to string type
    df_kafka_string=df_kafka.selectExpr("CAST(value AS STRING) as value")

    # schema to read json format data
    ts_schema = StructType() \
        .add("id_str", StringType()) \
        .add("created_at", StringType()) \
        .add("text", StringType())

    #parse json data
    df_kafka_string_parsed=df_kafka_string.select(from_json(col("value"),ts_schema).alias("twts"))

    df_kafka_string_parsed_format=df_kafka_string_parsed.select("twts.*")
    df_kafka_string_parsed_format.printSchema()


    df=df_kafka_string_parsed_format.writeStream \
    .trigger(processingTime="1 seconds") \
    .outputMode("update")\
    .option("truncate","false")\
    .format("console")\
    .start()

    df.awaitTermination()

【问题讨论】:

  • 请提供您的代码。
  • @AchyutVyas 我已经编辑了我的问题并提供了代码。请查看。
  • 如果删除 kafka-clients JAR 会发生什么?此外,Windows 对文件路径使用双反斜杠,而不是 //(正斜杠不需要“转义”)

标签: apache-spark pyspark apache-kafka spark-kafka-integration


【解决方案1】:

错误(NoClassDefFound,后跟kafka010 包)表示spark-sql-kafka-0-10 缺少对org.apache.commons:commons-pool2:2.6.2 的传递依赖,正如您可以see here 一样

您也可以下载该 JAR,或者您可以更改代码以使用 --packages 而不是 spark.jars 选项,并让 Ivy 处理下载传递依赖项

import os
os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache...'

spark = SparkSession.bulider...

【讨论】:

  • 你能扩展一下吗?我通过 sparkSesh = SparkSession.builder.config("spark.driver.extraClassPath", "/home/ubuntu/jars/spark-sql-kafka-0-10_2.12-3.1.2.jar,/home 添加了 .jar /ubuntu/jars/commons-pool2-2.11.0.jar") - --packages 替代方案是什么?
  • @James 不用下载 JAR,而是使用 Maven 坐标(以冒号分隔的格式)mvnrepository.com/artifact/org.apache.spark/…
猜你喜欢
  • 1970-01-01
  • 2022-10-26
  • 2016-06-11
  • 2011-01-21
  • 2023-03-30
  • 1970-01-01
  • 2018-04-04
  • 1970-01-01
  • 2013-02-19
相关资源
最近更新 更多