You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 00:34:55