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

Apache Flink运行Jar作业报错:找不到main(String[])方法

问题描述

通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 21:35:07