Spark独立模式运行含Kafka的Jar包时出现NoClassDefFoundError错误
问题描述
我在Spark中加载Jar包时遇到了以下错误:
Exception in thread "main" java.lang.NoClassDefFoundError: org/apache/kafka/clients/producer/Producer at com.sabd.project1.spark.Queries.calculateQuery1(Queries.java:53) at com.sabd.project1.Start.main(Start.java:54) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52) at org.apache.spark.deploy.SparkSubmit$.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:879) at org.apache.spark.deploy.SparkSubmit$.doRunMain$1(SparkSubmit.scala:197) at org.apache.spark.deploy.SparkSubmit$.submit(SparkSubmit.scala:227) at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:136) at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala) Caused by: java.lang.ClassNotFoundException: org.apache.kafka.clients.producer.Producer at java.net.URLClassLoader.findClass(URLClassLoader.java:381) at java.lang.ClassLoader.loadClass(ClassLoader.java:424) at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
我用Maven编译Jar包,对应的pom.xml配置如下:
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>SABD-Cenciarelli_Talone</groupId> <artifactId>Project1</artifactId> <version>1</version> <properties> <maven.compiler.target>1.8</maven.compiler.target> <maven.compiler.source>1.8</maven.compiler.source> </properties> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>3.6.1</version> <configuration> <source>1.8</source> <target>1.8</target> </configuration> </plugin> </plugins> </build> <dependencies> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.11</artifactId> <version>2.3.0</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.11</artifactId> <version>2.3.0</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka_2.11</artifactId> <version>1.1.0</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>1.1.0</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams</artifactId> <version>1.1.0</version> </dependency> <dependency> <groupId>com.google.code.gson</groupId> <artifactId>gson</artifactId> <version>2.8.2</version> </dependency> <dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk</artifactId> <version>1.11.318</version> </dependency> <dependency> <groupId>com.googlecode.json-simple</groupId> <artifactId>json-simple</artifactId> <version>1.1.1</version> </dependency> </dependencies> </project>
当用Spark本地模式(.setMaster("local"))运行时,错误就消失了。我用mvn clean install编译Jar包,怀疑是Maven没把Kafka相关类嵌入到Jar包里,请问该怎么解决?
解决方案
你的怀疑完全没错——本地模式下,Maven会自动把项目依赖加载到类路径里,所以程序能正常运行;但用spark-submit提交任务时,默认的Maven打包只会包含你自己写的代码,不会把依赖的第三方类(比如Kafka的客户端类)一起打包进去,这就导致Spark找不到这些类,抛出ClassNotFoundException。这里有两种非常实用的解决办法:
方法1:用Maven Shade插件打包成Fat Jar
Shade插件的作用就是把所有依赖的类都合并到一个Jar包中(也就是大家常说的Fat Jar),这样Spark提交任务时,只需要这一个Jar包就能找到所有需要的类。你只需要修改pom.xml的<build>部分,添加Shade插件配置:
<build> <plugins> <!-- 保留原有的编译插件 --> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>3.6.1</version> <configuration> <source>1.8</source> <target>1.8</target> </configuration> </plugin> <!-- 添加Shade打包插件 --> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.2.4</version> <executions> <execution> <phase>package</phase> <goals> <goal>shade</goal> </goals> <configuration> <!-- 指定主类,运行Jar时无需手动指定 --> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>com.sabd.project1.Start</mainClass> </transformer> </transformers> <!-- 排除签名文件,避免打包冲突 --> <filters> <filter> <artifact>*:*</artifact> <excludes> <exclude>META-INF/*.SF</exclude> <exclude>META-INF/*.DSA</exclude> <exclude>META-INF/*.RSA</exclude> </excludes> </filter> </filters> </configuration> </execution> </executions> </plugin> </plugins> </build>
修改完成后,运行mvn clean package(或mvn clean install),就会在target目录下生成一个包含所有依赖的Fat Jar(文件名通常为Project1-1-shaded.jar或类似格式),用这个Jar包提交Spark任务即可。
方法2:在spark-submit时指定依赖包
如果你不想打包成Fat Jar,也可以在提交Spark任务时,用--packages参数直接指定需要的Kafka依赖,Spark会自动从Maven仓库下载这些依赖并添加到类路径中:
spark-submit \ --class com.sabd.project1.Start \ --packages org.apache.kafka:kafka_2.11:1.1.0,org.apache.kafka:kafka-clients:1.1.0,org.apache.kafka:kafka-streams:1.1.0 \ your-jar-file.jar
不过这种方法每次提交任务时,若本地无缓存则会重新下载依赖,速度可能较慢,且需要确保集群能正常访问Maven仓库。
额外优化建议
对于Spark核心依赖(比如spark-core、spark-sql),建议在pom.xml中添加<scope>provided</scope>,因为Spark集群环境已经自带这些类,打包进去反而可能引发版本冲突:
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.11</artifactId> <version>2.3.0</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.11</artifactId> <version>2.3.0</version> <scope>provided</scope> </dependency>
这样Shade插件就不会把这些依赖打包进去,既能减小Jar包体积,也能避免不必要的冲突。
内容的提问来源于stack exchange,提问作者andrea cenciarelli

