Spark Streaming调用from_avro报NoClassDefFoundError求助
问题:Spark读取Kafka Avro序列化数据时
from_avro函数抛出类未找到异常 读取Kafka流中字符串类型Key正常,但调用from_avro反序列化Avro类型Value时抛出以下异常:
Exception in thread "main" java.lang.NoClassDefFoundError: org/apache/spark/sql/avro/functions$ at DataFrameExample$.main(DataFrameExample.scala:41) at DataFrameExample.main(DataFrameExample.scala) at java.base/jdk.internal.reflect.DirectMethodHandleAccessor.invoke(DirectMethodHandleAccessor.java:104) at java.base/java.lang.reflect.Method.invoke(Method.java:578) at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52) at org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:958) at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:180) at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:203) at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:90) at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:1046) at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:1055) at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala) Caused by: java.lang.ClassNotFoundException: org.apache.spark.sql.avro.functions$ at java.base/java.net.URLClassLoader.findClass(URLClassLoader.java:445) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:588) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:521) ... 12 more
当前SBT配置
name := "test" organization := "com.databricks" version := "1" scalaVersion := "2.12.17" // Spark Information val sparkVersion = "3.3.0" // allows us to include spark packages resolvers += "bintray-spark-packages" at "https://dl.bintray.com/spark-packages/maven/" resolvers += "Typesafe Simple Repository" at "https://repo.typesafe.com/typesafe/simple/maven-releases/" resolvers += "MavenRepository" at "https://mvnrepository.com/" libraryDependencies ++= Seq( "org.apache.spark" %% "spark-core" % sparkVersion % Provided, "org.apache.spark" %% "spark-streaming" % sparkVersion % Provided, "org.apache.spark" %% "spark-sql" % sparkVersion % Provided, "org.apache.spark" %% "spark-avro" % sparkVersion % Provided )
Scala代码
import org.apache.spark.sql._ import org.apache.spark.sql.types._ import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.avro.functions._ import java.nio.file.{Files, Paths} object DataFrameExample extends Serializable { def main(args: Array[String]) = { val spark = SparkSession .builder() .appName("Spark Example") .config("spark.sql.warehouse.dir", "/user/hive/warehouse") .getOrCreate() import spark.implicits._ spark.sparkContext.setLogLevel("WARN"); val currentDirectory = new java.io.File(".").getCanonicalPath println(currentDirectory); val df = spark.readStream .format("kafka") .option( "kafka.bootstrap.servers", "localhost:9091,localhost:9093,localhost:9094" ) .option("startingOffsets", "latest") .option("subscribe", "test-person-activity-partitions-replication-qwe") .load() val jsonFormatSchema = new String( Files.readAllBytes( Paths.get("./src/main/resources/avro/person-activity.avsc") ) ) val df2 = df.select( df.col("key").cast("string"), from_avro($"value", jsonFormatSchema).as("value") ) df2.writeStream .format("console") .outputMode("append") .start() .awaitTermination() } }
运行命令
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 --class DataFrameExample --master local target/scala-2.12/test_2.12-1.jar
解决方案
1. 修复spark-avro依赖范围
spark-avro是Spark的非核心扩展模块,默认Spark集群环境不会自带该依赖。你当前将其标记为Provided,意味着编译时依赖存在,但运行时不会打包到Jar中,导致找不到类。
修改SBT中的spark-avro依赖,移除Provided:
libraryDependencies ++= Seq( // 其他依赖保持不变 "org.apache.spark" %% "spark-avro" % sparkVersion )
2. 或在spark-submit时指定spark-avro依赖
如果不想修改SBT配置,也可以在提交命令中通过--packages参数引入spark-avro依赖,与Kafka依赖一起指定:
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0,org.apache.spark:spark-avro_2.12:3.3.0 --class DataFrameExample --master local target/scala-2.12/test_2.12-1.jar
3. 关于spark-sql-kafka-0-10依赖的冲突问题
你提到在SBT中引入该依赖会报错,通常是因为依赖版本冲突。可以将其标记为Provided,编译时获取依赖用于代码提示,运行时通过spark-submit的--packages引入,避免冲突:
libraryDependencies ++= Seq( // 其他依赖保持不变 "org.apache.spark" %% "spark-sql-kafka-0-10" % sparkVersion % Provided )
内容的提问来源于stack exchange,提问作者NikolaRoncevic
相关产品推荐
相关产品推荐

