Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Développer un job JAR

Dernière mise à jour :Aug 19, 2026

L'API DataStream de Flink propose un modèle de programmation flexible permettant de créer des transformations, des opérations et des opérateurs personnalisés pour gérer une logique métier complexe et le traitement des données.

Compatibilité Apache Flink

L'API DataStream prise en charge par Realtime Compute for Apache Flink est entièrement compatible avec la version open source d'Apache Flink. Pour plus d'informations, consultez les pages Qu'est-ce qu'Apache Flink ? et Guide de programmation de l'API Flink DataStream.

Prérequis

  • Un environnement de développement intégré (IDE), tel qu'IntelliJ IDEA, est installé.

  • Maven 3.6.3 ou une version ultérieure est installé.

  • Le développement de jobs prend uniquement en charge JDK 8 et JDK 11.

  • Développez votre job JAR localement avant de le déployer et de l'exécuter dans la console Realtime Compute for Apache Flink.

Avant de commencer

Préparez les sources de données requises à l'avance, car cet exemple utilise des connecteurs de source de données.

Remarque
  • Cet exemple utilise ApsaraMQ for Kafka (2.6.2) et ApsaraDB RDS for MySQL (8,0) comme sources de données.

  • Si vous disposez de sources de données gérées par vos soins nécessitant un accès au réseau public ou un accès inter-VPC, consultez la section Options de connectivité réseau.

  • Si vous ne disposez pas de source de données ApsaraMQ for Kafka, achetez et déployez une instance. Pour plus d'informations, consultez la section Étape 2 : Acheter et déployer une instance. Lors du déploiement de l'instance, assurez-vous qu'elle se trouve dans le même VPC que votre espace de travail Realtime Compute for Apache Flink.

  • Si vous ne disposez pas de source de données ApsaraDB RDS for MySQL, achetez une instance ApsaraDB RDS for MySQL. Pour plus d'informations, consultez la section Étape 1 : Créer une instance ApsaraDB RDS for MySQL et configurer une base de données. Lors de l'achat de l'instance, assurez-vous qu'elle se trouve dans la même région et le même VPC que votre espace de travail Realtime Compute for Apache Flink.

Développer le job

Configurer les dépendances de l'environnement Flink

Remarque

Pour éviter les conflits de dépendances JAR, suivez ces directives :

  • ${flink.version} spécifie la version de Flink pour l'exécution du job. Cette version doit être cohérente avec la version Flink du moteur VVR que vous avez sélectionnée sur la page de déploiement du job. Par exemple, si vous sélectionnez le moteur vvr-8.0.9-flink-1.17 sur la page de déploiement du job, sa version Flink correspondante est 1.17.2. Pour plus d'informations sur les versions du moteur VVR, consultez la section Comment vérifier la version Flink du job actuel ?.

  • Pour les dépendances Flink, définissez la portée (scope) sur provided en ajoutant <scope>provided</scope>. Cela concerne principalement les dépendances hors connecteur du groupe org.apache.flink qui commencent par flink-.

  • Dans le code source de Flink, seules les méthodes explicitement annotées avec @Public ou @PublicEvolving sont des API publiques. Realtime Compute for Apache Flink garantit la compatibilité uniquement pour ces méthodes.

  • Si vous utilisez l'API DataStream d'un connecteur Flink intégré, utilisez les dépendances fournies (provided).

Voici les dépendances Flink de base. Vous devrez peut-être également ajouter des dépendances de journalisation. Pour une liste complète des dépendances, consultez l'Exemple de code complet à la fin de cette rubrique.

Dépendances Flink

         <!-- Apache Flink dependencies -->
        <!-- These dependencies are set to 'provided' because they should not be packaged into the JAR file. -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-java</artifactId>
            <version>${flink.version}</version>
            <scope>provided</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java</artifactId>
            <version>${flink.version}</version>
            <scope>provided</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients</artifactId>
            <version>${flink.version}</version>
            <scope>provided</scope>
        </dependency>

Dépendances et utilisation des connecteurs

Pour lire et écrire des données avec l'API DataStream, utilisez un connecteur DataStream pour vous connecter à Realtime Compute for Apache Flink. Le dépôt central Maven met à disposition des connecteurs VVR DataStream que vous pouvez utiliser pendant le développement.

Important

Utilisez uniquement les connecteurs spécifiés comme prenant en charge l'API DataStream dans la section Connecteurs pris en charge. N'utilisez pas un connecteur s'il n'est pas explicitement marqué comme compatible avec l'API DataStream, car ses interfaces et paramètres pourraient changer à l'avenir.

Vous pouvez utiliser un connecteur de l'une des manières suivantes :

(Recommandé) Télécharger en tant que dépendance supplémentaire

  1. Dans le fichier pom.xml de votre job, ajoutez le connecteur requis en tant que dépendance de projet avec la portée « provided ». Pour le fichier de dépendances complet, consultez l'Exemple de code complet à la fin de cette rubrique.

    Remarque
    • ${vvr.version} correspond à la version du moteur d'exécution du job. Par exemple, si votre job s'exécute sur le moteur vvr-8.0.9-flink-1.17 , la version Flink correspondante est 1.17.2 . Nous vous recommandons d'utiliser la dernière version du moteur. Pour plus d'informations sur les versions spécifiques, consultez la section Moteur.

    • Étant donné que le package JAR du connecteur est ajouté en tant que dépendance supplémentaire, il n'a pas besoin d'être inclus dans le JAR de l'application. Par conséquent, vous devez déclarer sa portée comme provided .

            <!-- Kafka connector dependency -->
            <dependency>
                <groupId>com.alibaba.ververica</groupId>
                <artifactId>ververica-connector-kafka</artifactId>
                <version>${vvr.version}</version>
                <scope>provided</scope>
            </dependency>
            <!-- MySQL connector dependency -->
            <dependency>
                <groupId>com.alibaba.ververica</groupId>
                <artifactId>ververica-connector-mysql</artifactId>
                <version>${vvr.version}</version>
                <scope>provided</scope>
            </dependency>
  2. Si vous devez développer de nouveaux connecteurs ou étendre les fonctionnalités des connecteurs existants, votre projet doit également dépendre des packages de connecteurs communs flink-connector-base ou ververica-connector-common .

            <!-- Basic dependency for Flink connector public interfaces -->
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-connector-base</artifactId>
                <version>${flink.version}</version>
            </dependency>
            <!-- Basic dependency for Alibaba Cloud connector public interfaces -->
            <dependency>
                <groupId>com.alibaba.ververica</groupId>
                <artifactId>ververica-connector-common</artifactId>
                <version>${vvr.version}</version>
            </dependency>
  3. Pour les configurations de connexion DataStream et les exemples de code, consultez la documentation du connecteur DataStream correspondant.

    Pour obtenir la liste des connecteurs prenant en charge l'API DataStream, consultez la section Connecteurs pris en charge.

  4. Déployez le job et ajoutez les packages JAR de connecteur correspondants dans le champ Additional Dependencies . Pour plus d'informations, consultez la section Déployer un job JAR . Vous pouvez télécharger vos propres connecteurs ou ceux fournis par Realtime Compute for Apache Flink. Pour télécharger le connecteur, consultez la page Connecteurs .

    Par exemple, téléchargez les packages JAR de connecteur ververica-connector-mysql-1.17-vvr-8.0.4-1.jar et ververica-connector-kafka-1.17-vvr-8.0.4-1.jar , ainsi que le fichier de configuration config.properties .

Empaqueter dans le JAR du job

  1. Ajoutez les connecteurs requis en tant que dépendances de projet dans le fichier pom.xml de votre job. Le code suivant montre un exemple pour les connecteurs Kafka et MySQL.

    Remarque
    • ${vvr.version} correspond à la version du moteur de l'environnement d'exécution du job. Par exemple, si votre job s'exécute sur la version du moteur vvr-8.0.9-flink-1.17 , sa version Flink correspondante est 1.17.2 . Nous vous recommandons d'utiliser la dernière version du moteur. Pour plus d'informations, consultez la section Moteurs .

    • Si vous empaquetez les connecteurs dans le JAR du job en tant que dépendances de projet, ils doivent utiliser la portée par défaut (compile).

            <!-- Kafka connector dependency -->
            <dependency>
                <groupId>com.alibaba.ververica</groupId>
                <artifactId>ververica-connector-kafka</artifactId>
                <version>${vvr.version}</version>
            </dependency>
            <!-- MySQL connector dependency -->
            <dependency>
                <groupId>com.alibaba.ververica</groupId>
                <artifactId>ververica-connector-mysql</artifactId>
                <version>${vvr.version}</version>
            </dependency>
  2. Si vous devez développer de nouveaux connecteurs ou étendre les fonctionnalités des connecteurs existants, votre projet nécessite également le package de connecteur commun flink-connector-base ou ververica-connector-common .

            <!-- Basic dependency for Flink connector public interfaces -->
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-connector-base</artifactId>
                <version>${flink.version}</version>
            </dependency>
            <!-- Basic dependency for Alibaba Cloud connector public interfaces -->
            <dependency>
                <groupId>com.alibaba.ververica</groupId>
                <artifactId>ververica-connector-common</artifactId>
                <version>${vvr.version}</version>
            </dependency>
  3. Pour les configurations de connexion DataStream et les exemples de code, consultez la documentation du connecteur DataStream correspondant.

    Pour obtenir la liste des connecteurs prenant en charge l'API DataStream, consultez la section Connecteurs pris en charge .

Lire des dépendances supplémentaires depuis OSS

Les jobs JAR Flink ne peuvent pas lire les fichiers de configuration locaux depuis la méthode main . À la place, téléchargez le fichier de configuration dans le bucket OSS de votre espace de travail, ajoutez-le en tant que dépendance supplémentaire lors du déploiement et lisez-le au moment de l'exécution. La section suivante fournit un exemple.

  1. Créez un fichier de configuration nommé config.properties pour éviter d'encoder en dur les identifiants dans votre code.

    # Kafka 
    bootstrapServers=host1:9092,host2:9092,host3:9092
    inputTopic=topic
    groupId=groupId
    # MySQL
    database.url=jdbc:mysql://localhost:3306/my_database
    database.username=username
    database.password=password
  2. Dans votre job JAR, utilisez du code pour lire le fichier config.properties stocké dans le bucket OSS.

    Méthode 1 : Lire depuis le bucket de l'espace de travail

    1. Dans le volet de navigation de gauche de la console Realtime Compute for Apache Flink, accédez à la page Artifacts et téléchargez le fichier.

    2. Au moment de l'exécution, les fichiers ajoutés dans le champ Additional Dependencies sont chargés dans le répertoire /flink/usrlib du pod où le job s'exécute.

    3. Le code suivant illustre la lecture de ce fichier de configuration.

                  Properties properties = new Properties();
                  Map<String,String> configMap = new HashMap<>();
                  try (InputStream input = new FileInputStream("/flink/usrlib/config.properties")) {
                      // Load the property file.
                      properties.load(input);
                      // Obtain the property values.
                      configMap.put("bootstrapServers",properties.getProperty("bootstrapServers")) ;
                      configMap.put("inputTopic",properties.getProperty("inputTopic"));
                      configMap.put("groupId",properties.getProperty("groupId"));
                      configMap.put("url",properties.getProperty("database.url")) ;
                      configMap.put("username",properties.getProperty("database.username"));
                      configMap.put("password",properties.getProperty("database.password"));
                  } catch (IOException ex) {
                      ex.printStackTrace();
                  }

    Méthode 2 : Lire depuis un bucket autorisé

    1. Téléchargez le fichier de configuration dans le bucket OSS cible.

    2. Utilisez OSSClient pour lire le fichier directement depuis OSS. Pour plus d'informations, consultez les sections Téléchargement en flux continu et Gérer les identifiants d'accès . Le code suivant fournit un exemple.

      OSS ossClient = new OSSClientBuilder().build("Endpoint", "AccessKeyId", "AccessKeySecret");
      try (OSSObject ossObject = ossClient.getObject("examplebucket", "exampledir/config.properties");
           BufferedReader reader = new BufferedReader(new InputStreamReader(ossObject.getObjectContent()))) {
          // read file and process ...
      } finally {
          if (ossClient != null) {
              ossClient.shutdown();
          }
      }

Écrire la logique métier

  1. Vous pouvez intégrer des sources de données externes dans les programmes de flux de données Flink. Un Watermark est une stratégie de calcul Flink basée sur la sémantique temporelle et souvent utilisée avec des horodatages. Par conséquent, cet exemple n'utilise pas de stratégie de watermark. Pour plus d'informations, consultez la page Stratégie de watermark .

             // Integrate the external data source into the Flink data stream program.
            // WatermarkStrategy.noWatermarks() indicates that no watermark strategy is used.
            DataStreamSource<String> stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka Source");
  2. La transformation d'opérateur convertit un DataStream<String> en un DataStream<Student> dans cet exemple. Pour des transformations d'opérateur et des méthodes de traitement plus complexes, consultez la page Opérateurs Flink .

              // Operator that converts the data structure to Student.
              DataStream<Student> source = stream
                    .map(new MapFunction<String, Student>() {
                        @Override
                        public Student map(String s) throws Exception {
                            // Data is separated by commas.
                            String[] data = s.split(",");
                            return new Student(Integer.parseInt(data[0]), data[1], Integer.parseInt(data[2]));
                        }
                    }).filter(student -> student.score >=60); // Filter for data where the score is 60 or higher.

Empaqueter le job

Empaquetez le job à l'aide du plugin maven-shade-plugin.

Important
  • Si vous ajoutez le connecteur en tant que dépendance supplémentaire, assurez-vous que la portée des dépendances du connecteur est définie sur provided lors de l'empaquetage du job.

  • Si vous empaquetez les connecteurs dans le JAR du job, utilisez la portée par défaut (compile).

Référence pour les dépendances maven-shade-plugin

<build>
        <plugins>
            <!-- Java Compiler -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <version>3.11.0</version>
                <configuration>
                    <source>${target.java.version}</source>
                    <target>${target.java.version}</target>
                </configuration>
            </plugin>
            <!-- We use the maven-shade-plugin to create a fat JAR that contains all necessary dependencies. -->
            <!-- Modify the value of <mainClass> if your program's entry point changes. -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-shade-plugin</artifactId>
                <version>3.2.0</version>
                <executions>
                    <execution>
                        <phase>package</phase>
                        <goals>
                            <goal>shade</goal>
                        </goals>
                        <!-- Exclude unnecessary dependencies. -->
                        <configuration>
                            <artifactSet>
                                <excludes>
                                    <exclude>org.apache.flink:force-shading</exclude>
                                    <exclude>com.google.code.findbugs:jsr305</exclude>
                                    <exclude>org.slf4j:*</exclude>
                                    <exclude>org.apache.logging.log4j:*</exclude>
                                </excludes>
                            </artifactSet>
                            <filters>
                                <filter>
                                    <!-- Do not copy the signatures from the META-INF directory.
                                    Otherwise, security exceptions may be thrown when you use the JAR file. -->
                                    <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>com.aliyun.FlinkDemo</mainClass>
                                </transformer>
                            </transformers>
                        </configuration>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>

Tester et déployer le job

  • Par défaut, Realtime Compute for Apache Flink ne peut pas accéder à Internet public, ce qui empêche les tests directs dans un environnement local. Nous vous recommandons d'effectuer des tests unitaires séparément. Pour plus d'informations, consultez la section Exécuter et déboguer un job avec des connecteurs localement .

  • Pour déployer le job JAR, consultez la section Déployer un job JAR .

    Remarque
    • Lors du déploiement, si vous utilisez des connecteurs en tant que dépendances supplémentaires, assurez-vous de télécharger les packages JAR pertinents.

    • Si vous devez lire un fichier de configuration, vous devez également le télécharger en tant que dépendance supplémentaire.

Exemple de code complet

Cet exemple traite les données provenant d'une source ApsaraMQ for Kafka et écrit le résultat dans un puits ApsaraDB RDS for MySQL. Cet exemple est fourni à titre de référence uniquement. Pour plus d'informations sur le style et la qualité du code, consultez le guide Style et qualité du code .

Remarque

Cet exemple omet les configurations des paramètres d'exécution tels que checkpoint, TTL et stratégie de redémarrage. Vous pouvez configurer ces paramètres sur la page deployment details après le déploiement du job. Étant donné que les paramètres définis dans le code ont une priorité plus élevée, nous vous recommandons de les configurer après le déploiement afin de simplifier les mises à jour futures. Pour plus d'informations, consultez la section Configurer les informations de déploiement du job .

FlinkDemo.java

package com.aliyun;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.connector.jdbc.JdbcStatementBuilder;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
import org.apache.flink.kafka.shaded.org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.io.FileInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
public class FlinkDemo {
    // Define the data structure.
    public static class Student {
        public int id;
        public String name;
        public int score;
        public Student(int id, String name, int score) {
            this.id = id;
            this.name = name;
            this.score = score;
        }
    }
    public static void main(String[] args) throws Exception {
        // Create a Flink execution environment.
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        Properties properties = new Properties();
        Map<String,String> configMap = new HashMap<>();
        try (InputStream input = new FileInputStream("/flink/usrlib/config.properties")) {
            // Load the property file.
            properties.load(input);
            // Obtain the property values.
            configMap.put("bootstrapServers",properties.getProperty("bootstrapServers")) ;
            configMap.put("inputTopic",properties.getProperty("inputTopic"));
            configMap.put("groupId",properties.getProperty("groupId"));
            configMap.put("url",properties.getProperty("database.url")) ;
            configMap.put("username",properties.getProperty("database.username"));
            configMap.put("password",properties.getProperty("database.password"));
        } catch (IOException ex) {
            ex.printStackTrace();
        }
        // Build Kafka source
        KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
                        .setBootstrapServers(configMap.get("bootstrapServers"))
                        .setTopics(configMap.get("inputTopic"))
                        .setStartingOffsets(OffsetsInitializer.latest())
                        .setGroupId(configMap.get("groupId"))
                        .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))
                        .build();
        // Integrate the external data source into the Flink data stream program.
        // WatermarkStrategy.noWatermarks() indicates that no watermark strategy is used.
        DataStreamSource<String> stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka Source");
        // Filter for scores of 60 or higher.
        DataStream<Student> source = stream
                .map(new MapFunction<String, Student>() {
                    @Override
                    public Student map(String s) throws Exception {
                        String[] data = s.split(",");
                        return new Student(Integer.parseInt(data[0]), data[1], Integer.parseInt(data[2]));
                    }
                }).filter(Student -> Student.score >=60);
        source.addSink(JdbcSink.sink("INSERT IGNORE INTO student (id, username, score) VALUES (?, ?, ?)",
                new JdbcStatementBuilder<Student>() {
                    public void accept(PreparedStatement ps, Student data) {
                        try {
                            ps.setInt(1, data.id);
                            ps.setString(2, data.name);
                            ps.setInt(3, data.score);
                        } catch (SQLException e) {
                            throw new RuntimeException(e);
                        }
                    }
                },
                new JdbcExecutionOptions.Builder()
                        .withBatchSize(5) // Number of records per batch write.
                        .withBatchIntervalMs(2000) // Batch interval in milliseconds.
                        .build(),
                new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
                        .withUrl(configMap.get("url"))
                        .withDriverName("com.mysql.cj.jdbc.Driver")
                        .withUsername(configMap.get("username"))
                        .withPassword(configMap.get("password"))
                        .build()
        )).name("Sink MySQL");
        env.execute("Flink Demo");
    }
}

pom.xml

<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/maven-v4_0_0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>com.aliyun</groupId>
    <artifactId>FlinkDemo</artifactId>
    <version>1.0-SNAPSHOT</version>
    <name>FlinkDemo</name>
    <packaging>jar</packaging>
    <properties>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
        <flink.version>1.17.1</flink.version>
        <vvr.version>1.17-vvr-8.0.4-1</vvr.version>
        <target.java.version>1.8</target.java.version>
        <maven.compiler.source>${target.java.version}</maven.compiler.source>
        <maven.compiler.target>${target.java.version}</maven.compiler.target>
        <log4j.version>2.14.1</log4j.version>
    </properties>
    <dependencies>
        <!-- Apache Flink dependencies -->
        <!-- These dependencies are set to 'provided' because they should not be packaged into the JAR file. -->
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-java</artifactId>
            <version>${flink.version}</version>
            <scope>provided</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java</artifactId>
            <version>${flink.version}</version>
            <scope>provided</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients</artifactId>
            <version>${flink.version}</version>
            <scope>provided</scope>
        </dependency>
        <!-- Add connector dependencies here. They must be in the default scope (compile). -->
        <dependency>
            <groupId>com.alibaba.ververica</groupId>
            <artifactId>ververica-connector-kafka</artifactId>
            <version>${vvr.version}</version>
        </dependency>
        <dependency>
            <groupId>com.alibaba.ververica</groupId>
            <artifactId>ververica-connector-mysql</artifactId>
            <version>${vvr.version}</version>
        </dependency>
        <!-- Add a logging framework to generate console output at runtime. -->
        <!-- By default, these dependencies are excluded from the application JAR. -->
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-slf4j-impl</artifactId>
            <version>${log4j.version}</version>
            <scope>runtime</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-api</artifactId>
            <version>${log4j.version}</version>
            <scope>runtime</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.logging.log4j</groupId>
            <artifactId>log4j-core</artifactId>
            <version>${log4j.version}</version>
            <scope>runtime</scope>
        </dependency>
    </dependencies>
    <build>
        <plugins>
            <!-- Java Compiler -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <version>3.11.0</version>
                <configuration>
                    <source>${target.java.version}</source>
                    <target>${target.java.version}</target>
                </configuration>
            </plugin>
            <!-- We use the maven-shade-plugin to create a fat JAR that contains all necessary dependencies. -->
            <!-- Modify the value of <mainClass> if your program's entry point changes. -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-shade-plugin</artifactId>
                <version>3.2.0</version>
                <executions>
                    <execution>
                        <phase>package</phase>
                        <goals>
                            <goal>shade</goal>
                        </goals>
                        <!-- Exclude unnecessary dependencies. -->
                        <configuration>
                            <artifactSet>
                                <excludes>
                                    <exclude>org.apache.flink:force-shading</exclude>
                                    <exclude>com.google.code.findbugs:jsr305</exclude>
                                    <exclude>org.slf4j:*</exclude>
                                    <exclude>org.apache.logging.log4j:*</exclude>
                                </excludes>
                            </artifactSet>
                            <filters>
                                <filter>
                                    <!-- Do not copy the signatures from the META-INF directory.
                                    Otherwise, security exceptions may be thrown when you use the JAR file. -->
                                    <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>com.aliyun.FlinkDemo</mainClass>
                                </transformer>
                            </transformers>
                        </configuration>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>

Documentation connexe