Apache Flink运行Jar作业报错:找不到main(String[])方法
Apache Flink 运行Jar包报错问题排查
问题描述
通过Eclipse的maven install功能和mvn clean package命令(借助maven-assembly-plugin)构建Flink项目Jar包后,在Flink主目录执行以下命令时触发异常:
./bin/flink run -c tbdm.kafka_flink_integration.Flink_Kafka_Receiver_Sender ./flink-kafka.jar
程序在Eclipse中可正常运行,但本地Ubuntu机器上的Flink集群(可通过localhost:8081访问GUI控制台)既无法通过命令行运行该Jar,也无法通过控制台加载Jar。
报错信息
java.lang.RuntimeException: Could not look up the main(String[]) method from the class tbdm.kafka_flink_integration.Flink_Kafka_Receiver_Sender: org/apache/flink/streaming/connectors/kafka/FlinkKafkaProducer at org.apache.flink.client.program.PackagedProgram.hasMainMethod(PackagedProgram.java:315) and so on...
环境与代码
- Flink版本:1.16.1,所有Flink依赖均为对应版本
- 开发语言:Java
- 主类功能:从Kafka Topic读取数据,转为大写后存入MongoDB,同时输出到另一个Kafka Topic
主类代码
package tbdm.kafka_flink_integration; import java.util.Properties; import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer; import com.mongodb.client.MongoClient; import com.mongodb.client.MongoCollection; import com.mongodb.client.MongoDatabase; import com.mongodb.client.MongoClients; import org.bson.Document; public class Flink_Kafka_Receiver_Sender { public static void main(String[] args) throws Exception { String inputTopic = "flink_input"; String outputTopic = "flink_output"; String server = "localhost:9092"; StreamStringOperation(server, inputTopic, outputTopic); } public static void StreamStringOperation(String server, String inputTopic, String outputTopic) throws Exception { StreamExecutionEnvironment environment = StreamExecutionEnvironment.getExecutionEnvironment(); Properties props = new Properties(); props.setProperty("bootstrap.servers", server); FlinkKafkaConsumer<String> flinkKafkaConsumer = new FlinkKafkaConsumer<>(inputTopic, new SimpleStringSchema(), props); FlinkKafkaProducer<String> flinkKafkaProducer = new FlinkKafkaProducer<>(server, outputTopic, new SimpleStringSchema()); DataStream<String> stringInputStream = environment.addSource(flinkKafkaConsumer); stringInputStream.map(new StringCapitalizer()).addSink(flinkKafkaProducer); environment.execute(); } public static class StringCapitalizer implements MapFunction<String, String> { @Override public String map(String data) throws Exception { System.out.println(data.toUpperCase()); // 创建MongoDB连接 MongoClient mongoClient = MongoClients.create("mongodb://localhost:27017"); MongoDatabase database = mongoClient.getDatabase("myDatabase"); // 选择目标集合 MongoCollection<Document> collection = database.getCollection("myCollection"); // 构建并插入文档 Document document = new Document(); document.append("stringa", data.toUpperCase()); collection.insertOne(document); return data.toUpperCase(); } } public static void StreamConsumer(String inputTopic, String server) throws Exception { StreamExecutionEnvironment environment = StreamExecutionEnvironment.getExecutionEnvironment(); Properties props = new Properties(); props.setProperty("bootstrap.servers", server); DataStream<String> stringInputStream = environment.addSource(new FlinkKafkaConsumer<String>(inputTopic, new SimpleStringSchema(), props)); stringInputStream.map(new MapFunction<String, String>() { private static final long serialVersionUID = -999736771747691234L; public String map(String value) throws Exception { return "Receiving from Kafka : " + value; } }).print(); environment.execute(); } public static FlinkKafkaProducer<String> createStringProducer(String output_topic, String kafkaAddress) { Properties props = new Properties(); props.setProperty("bootstrap.servers", kafkaAddress); return new FlinkKafkaProducer<String>(output_topic, new SimpleStringSchema(), props); } }
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>tbdm</groupId> <artifactId>kafka-flink-integration</artifactId> <version>0.0.1-SNAPSHOT</version> <name>kafka-flink-integration</name> <url>http://www.example.com</url> <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <maven.compiler.source>1.7</maven.compiler.source> <maven.compiler.target>1.7</maven.compiler.target> </properties> <dependencies> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>4.11</version> <scope>test</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-core</artifactId> <version>1.16.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>1.16.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.16.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients</artifactId> <version>1.16.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.mongodb</groupId> <artifactId>mongo-java-driver</artifactId> <version>3.12.12</version> </dependency> </dependencies> <build> <pluginManagement> <plugins> <plugin> <artifactId>maven-clean-plugin</artifactId> <version>3.1.0</version> </plugin> <plugin> <artifactId>maven-resources-plugin</artifactId> <version>3.0.2</version> </plugin> <plugin> <artifactId>maven-compiler-plugin</artifactId> <version>3.8.0</version> </plugin> <plugin> <artifactId>maven-surefire-plugin</artifactId> <version>2.22.1</version> </plugin> <plugin> <artifactId>maven-jar-plugin</artifactId> <version>3.0.2</version> </plugin> <plugin> <artifactId>maven-install-plugin</artifactId> <version>2.5.2</version> </plugin> <plugin> <artifactId>maven-deploy-plugin</artifactId> <version>2.8.2</version> </plugin> <plugin> <artifactId>maven-site-plugin</artifactId> <version>3.7.1</version> </plugin> <plugin> <artifactId>maven-project-info-reports-plugin</artifactId> <version>3.0.0</version> </plugin> <plugin> <artifactId>maven-assembly-plugin</artifactId> <version>3.5.0</version> <configuration> <archive> <manifest> <mainClass>tbdm.kafka_flink_integration.Flink_Kafka_Receiver_Sender</mainClass> </manifest> </archive> <descriptorRefs> <descriptorRef>jar-with-dependencies</descriptorRef> </descriptorRefs> </configuration> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> </plugin> </plugins> </pluginManagement> </build> </project>
解决方案
1. 修复依赖缺失问题
报错核心是找不到FlinkKafkaProducer类,原因是Flink集群默认不包含Kafka连接器依赖,且pom中该依赖的scope设为provided,打包时不会包含进Jar。解决方式二选一:
- 将
flink-connector-kafka的scope改为compile,重新打包后运行 - 手动将
flink-connector-kafka-1.16.1.jar复制到Flink安装目录的lib文件夹下,重启Flink集群
2. 修正maven-assembly-plugin配置
当前assembly插件放在pluginManagement节点下,仅用于版本管理,不会实际执行打包操作。需将其移到build/plugins节点下:
<build> <plugins> <plugin> <artifactId>maven-assembly-plugin</artifactId> <version>3.5.0</version> <configuration> <archive> <manifest> <mainClass>tbdm.kafka_flink_integration.Flink_Kafka_Receiver_Sender</mainClass> </manifest> </archive> <descriptorRefs> <descriptorRef>jar-with-dependencies</descriptorRef> </descriptorRefs> </configuration> <executions> <execution> <id>make-assembly</id> <phase>package</phase> <goals> <goal>single</goal> </goals> </execution> </executions> </plugin> </plugins> </build>
3. 代码优化建议
- 避免在
map函数中每次创建MongoDB连接,可将连接逻辑移至open方法中复用,或使用Flink官方MongoDB连接器(flink-connector-mongodb) - 检查Kafka的
bootstrap.servers配置,确保Flink集群能正常访问Kafka服务
内容的提问来源于stack exchange,提问作者niccolonator
相关产品推荐
相关产品推荐

