Flink 1.16.2作业执行报错ClassNotFoundException: KafkaSource求助
问题
使用Apache Flink 1.16.2版本运行Flink作业,执行命令:
./bin/flink run /Users/spartacus/icu-alarm/target/flink-kafka-stroke-risk-1.0-SNAPSHOT.jar
时遇到以下错误:
java.lang.NoClassDefFoundError: org/apache/flink/connector/kafka/source/KafkaSource at hes.cs63.CEPMonitor.StrokeRiskAlarm.main(StrokeRiskAlarm.java:30) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:566) at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:355) at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:222) at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:98) at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:843) at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:240) at org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1087) at org.apache.flink.client.cli.CliFrontend.lambda$main$10(CliFrontend.java:1165) at org.apache.flink.runtime.security.contexts.NoOpSecurityContext.runSecured(NoOpSecurityContext.java:28) at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1165) Caused by: java.lang.ClassNotFoundException: org.apache.flink.connector.kafka.source.KafkaSource at java.base/java.net.URLClassLoader.findClass(URLClassLoader.java:476) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:588) at org.apache.flink.util.FlinkUserCodeClassLoader.loadClassWithoutExceptionHandling(FlinkUserCodeClassLoader.java:67) at org.apache.flink.util.ChildFirstClassLoader.loadClassWithoutExceptionHandling(ChildFirstClassLoader.java:65) at org.apache.flink.util.FlinkUserCodeClassLoader.loadClass(FlinkUserCodeClassLoader.java:51) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:521) ... 14 more
配置Kafka作为数据源的相关代码:
// KafkaSource setup KafkaSource<String> source = KafkaSource.<String>builder() .setBootstrapServers("localhost:9092") .setGroupId("stroke-risk-group") .setTopics(Arrays.asList("patient-data-topic")) .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class)) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStreamSource<String> patientData = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");
额外信息:
- Flink版本: 1.16.2
- Java版本: 11
- IDE: IntelliJ IDEA
原因与解决方案
错误原因
该错误本质是Flink Kafka连接器的依赖未被正确打包到作业JAR中,导致运行时类路径找不到KafkaSource相关类。Flink核心依赖默认由集群提供,但Kafka连接器属于可选依赖,需手动配置打包策略确保相关类被包含在最终JAR内。
解决步骤
1. 检查Maven依赖配置
确保pom.xml中已添加Flink Kafka连接器依赖,且scope设置为compile(默认),不要设为provided(否则Maven打包时会排除该依赖):
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>1.16.2</version> </dependency> <!-- 需同时包含Flink Streaming API依赖 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.16.2</version> <scope>compile</scope> </dependency>
2. 构建包含所有依赖的Fat Jar
默认Maven打包仅包含自有代码,需使用maven-shade-plugin将所有依赖打包成一个Fat Jar。在pom.xml的<build><plugins>中添加以下配置:
<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> <artifactSet> <excludes> <!-- 排除Flink核心依赖,避免与集群冲突 --> <exclude>org.apache.flink:flink-shaded-*</exclude> <exclude>org.apache.flink:flink-core</exclude> <exclude>org.apache.flink:flink-java</exclude> <exclude>org.apache.flink:flink-streaming-java</exclude> </excludes> </artifactSet> <filters> <filter> <artifact>*:*</artifact> <excludes> <exclude>META-INF/*.SF</exclude> <exclude>META-INF/*.DSA</exclude> <exclude>META-INF/*.RSA</exclude> </excludes> </filter> </filters> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <!-- 指定作业主类的全路径 --> <mainClass>hes.cs63.CEPMonitor.StrokeRiskAlarm</mainClass> </transformer> </transformers> </configuration> </execution> </executions> </plugin>
3. 重新打包并运行
执行Maven打包命令:
mvn clean package
打包完成后,使用target目录下生成的Fat Jar(体积通常远大于原JAR)重新运行作业:
./bin/flink run /Users/spartacus/icu-alarm/target/flink-kafka-stroke-risk-1.0-SNAPSHOT.jar
4. 备选方案:运行时指定依赖路径
若不想打包Fat Jar,可在运行时通过-jars参数指定Kafka连接器JAR的路径,同时用-c指定主类:
./bin/flink run -c hes.cs63.CEPMonitor.StrokeRiskAlarm -jars /path/to/flink-connector-kafka-1.16.2.jar /Users/spartacus/icu-alarm/target/flink-kafka-stroke-risk-1.0-SNAPSHOT.jar
该方式适合快速测试,生产环境推荐使用Fat Jar避免依赖路径问题。
内容的提问来源于stack exchange,提问作者Ishaan Adarsh
相关产品推荐
相关产品推荐

