Tous les produits
Search
Centre de documentation

E-MapReduce:Diffuser des données Kafka vers Alibaba Cloud OSS avec Flink

Dernière mise à jour :Aug 09, 2026

Ce tutoriel vous explique comment exécuter un job de streaming Flink sur un cluster Dataflow E-MapReduce (EMR) qui lit depuis un topic Kafka et écrit dans Object Storage Service (OSS) avec une sémantique exactly-once. Le job utilise JindoFS — disponible dans les clusters Dataflow à partir d'EMR V3.37.1 — pour écrire dans OSS en utilisant des chemins oss://, de la même manière que vous écrivez dans le système de fichiers distribué Hadoop (HDFS).

Prérequis

Avant de commencer, assurez-vous d'avoir :

  • Activé EMR et OSS sur votre compte Alibaba Cloud

  • Accordé les permissions requises à votre compte Alibaba Cloud ou utilisateur RAM. Pour plus de détails, consultez la rubrique Attribuer des rôles

Fonctionnement

Le job Flink utilise StreamingFileSink avec la politique OnCheckpointRollingPolicy. Selon cette politique, les fichiers de sortie sont validés dans OSS uniquement lorsqu'un checkpoint est terminé, et non en continu. Avec un intervalle de checkpoint de 30 secondes, cela signifie que de nouveaux fichiers de sortie apparaissent environ toutes les 30 secondes. Étant donné qu'il s'agit d'un job de streaming, il s'exécute en continu et accumule les fichiers de sortie jusqu'à ce que vous l'arrêtiez.

JindoFS gère les écritures dans OSS. Si le bucket OSS et le cluster Dataflow appartiennent au même compte Alibaba Cloud, JindoFS lit et écrit dans le bucket en mode sans mot de passe ; aucune information d'identification n'est nécessaire dans le code du job.

Important

Étant donné que les fichiers de sortie s'accumulent à chaque checkpoint, arrêtez le job après avoir vérifié la sortie. Consultez la section « Arrêter le job » à l'étape 5.

Étape 1 : Préparer l'environnement

Remarque

Ce tutoriel utilise EMR V3.43.1.

  1. Créez un cluster Dataflow avec les services Flink et Kafka. Pour plus de détails, consultez la rubrique Créer un cluster.

  2. Créez un bucket OSS dans la même région que le cluster Dataflow. Pour plus de détails, consultez la rubrique Créer des buckets.

Étape 2 : Préparer le package JAR

L'exemple de code crée une source Kafka à l'aide de l'API source FLIP-27 et un puits OSS pris en charge par StreamingFileSink. Le chemin de sortie doit commencer par oss://. La fonctionnalité de checkpointing est activée avec le mode CheckpointingMode.EXACTLY_ONCE et un intervalle de 30 secondes ; c'est ce qui pilote la politique OnCheckpointRollingPolicy et détermine la fréquence de validation des fichiers.

public class OssDemoJob {

    public static void main(String[] args) throws Exception {
        ...

        // Check output oss dir
        Preconditions.checkArgument(
                params.get(OUTPUT_OSS_DIR).startsWith("oss://"),
                "outputOssDir should start with 'oss://'.");

        // Set up the streaming execution environment
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // Checkpoint is required
        env.enableCheckpointing(30000, CheckpointingMode.EXACTLY_ONCE);

        String outputPath = params.get(OUTPUT_OSS_DIR);

        // Build Kafka source with new Source API based on FLIP-27
        KafkaSource<Event> kafkaSource =
                KafkaSource.<Event>builder()
                        .setBootstrapServers(params.get(KAFKA_BROKERS_ARG))
                        .setTopics(params.get(INPUT_TOPIC_ARG))
                        .setStartingOffsets(OffsetsInitializer.latest())
                        .setGroupId(params.get(INPUT_TOPIC_GROUP_ARG))
                        .setDeserializer(new EventDeSerializationSchema())
                        .build();
        // DataStream Source
        DataStreamSource<Event> source =
                env.fromSource(
                        kafkaSource,
                        WatermarkStrategy.<Event>forMonotonousTimestamps()
                                .withTimestampAssigner((event, ts) -> event.getEventTime()),
                        "Kafka Source");

        StreamingFileSink<Event> sink =
                StreamingFileSink.forRowFormat(
                                new Path(outputPath), new SimpleStringEncoder<Event>("UTF-8"))
                        .withRollingPolicy(OnCheckpointRollingPolicy.build())
                        .build();
        source.addSink(sink);

        // Compile and submit the job
        env.execute();
    }
}
Remarque

Cet extrait présente la structure principale du programme. Modifiez-le selon vos besoins, par exemple en ajoutant un nom de package ou en ajustant l'intervalle de checkpoint, avant la compilation. Pour obtenir des instructions sur la création d'un package JAR, consultez la documentation officielle de Flink.

Téléchargez le code source complet depuis le dépôt dataflow-demo, ou utilisez directement le fichier précompilé dataflow-oss-demo-1.0-SNAPSHOT.jar.

Pour compiler le fichier JAR vous-même :

  1. Accédez au répertoire racine du projet téléchargé.

  2. Exécutez la commande suivante :

    mvn clean package

    Le fichier JAR généré, dataflow-oss-demo-1.0-SNAPSHOT.jar, est placé dans le répertoire dataflow-demo/dataflow-oss-demo/target/ en fonction de l'artifactId défini dans le fichier pom.xml.

Étape 3 : Créer un topic Kafka et générer des données

  1. Connectez-vous au cluster Dataflow via SSH. Pour plus de détails, consultez la rubrique Se connecter à un cluster.

  2. Créez un topic Kafka pour les tests :

    kafka-topics.sh --create  --bootstrap-server core-1-1:9092 \
        --replication-factor 2  \
        --partitions 3  \
        --topic kafka-test-topic

    En cas de succès, l'interface CLI renvoie :

    Created topic kafka-test-topic.
  3. Écrivez des données de test dans le topic.

    1. Ouvrez la console du producteur Kafka :

      kafka-console-producer.sh --broker-list core-1-1:9092 --topic  kafka-test-topic
    2. Saisissez les cinq enregistrements suivants :

      1,Ken,0,1,1662022777000
      1,Ken,0,2,1662022777000
      1,Ken,0,3,1662022777000
      1,Ken,0,4,1662022777000
      1,Ken,0,5,1662022777000
    3. Appuyez sur Ctrl+C pour quitter la console du producteur Kafka.

Étape 4 : Exécuter un job Flink

  1. Connectez-vous au cluster Dataflow via SSH. Pour plus de détails, consultez la rubrique Se connecter à un cluster.

  2. Téléchargez le fichier dataflow-oss-demo-1.0-SNAPSHOT.jar dans le répertoire racine du cluster Dataflow.

    Remarque

    Cet exemple télécharge le fichier dans le répertoire racine. Choisissez un autre répertoire si nécessaire.

  3. Soumettez le job Flink en mode Per-Job :

    flink run -t yarn-per-job -d -c com.alibaba.ververica.dataflow.demo.oss.OssDemoJob \
        /dataflow-oss-demo-1.0-SNAPSHOT.jar  \
        --outputOssDir oss://xung****-flink-dlf-test/oss_kafka_test \
        --kafkaBrokers core-1-1:9092 \
        --inputTopic kafka-test-topic \
        --inputTopicGroup my-group
    Paramètre Description Valeurs valides Valeur dans cet exemple
    --outputOssDir Chemin OSS vers lequel écrire les données. Doit commencer par oss://. Tout chemin OSS valide commençant par oss:// oss://xung****-flink-dlf-test/oss_kafka_test
    --kafkaBrokers Adresse du broker Kafka host:port core-1-1:9092
    --inputTopic Topic Kafka à lire Nom de tout topic existant kafka-test-topic
    --inputTopicGroup Groupe de consommateurs Kafka Nom valide de groupe de consommateurs my-group

    Pour connaître les autres modes de soumission, consultez la rubrique Utilisation de base.

    La sortie de journal suivante est renvoyée :

  4. Vérifiez que le job est en cours d'exécution :

    flink list -t yarn-per-job -Dyarn.application.id=<appId>

    Remplacez <appId> par l'ID d'application renvoyé après la soumission du job. Dans cet exemple, l'ID d'application est application_1670236019397_0003.

Étape 5 : Afficher la sortie

Le job valide les fichiers de sortie à chaque checkpoint (toutes les 30 secondes). Affichez la sortie, puis arrêtez le job.

Afficher la sortie dans la console OSS :

  1. Connectez-vous à la console OSS.

  2. Dans le volet de navigation de gauche, cliquez sur Buckets. Sur la page Buckets, cliquez sur le nom du bucket que vous avez créé.

  3. Sur la page Objects, accédez au répertoire de sortie que vous avez spécifié dans le paramètre --outputOssDir.

    Sur la page Objects, vous pouvez voir des répertoires de sortie nommés par horodatage. Accédez à un répertoire pour trouver les fichiers de résultats générés (par exemple, part-0-0).

Afficher la sortie depuis l'interface CLI :

Exécutez la commande suivante sur le cluster Dataflow :

hdfs dfs -cat oss://<YOUR_TARGET_BUCKET>/oss_kafka_test/<DATE_DIR>/part-0-0

La sortie d'exemple suivante est renvoyée :

[root@master-1-xxx ~]# hdfs dfs -cat oss://xxx2022-12-06--14/part-0-0
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,1,1662022777000
Sink to OSS: 1,Ken,0,1,1662022777000

Arrêter le job :

Après avoir vérifié la sortie, arrêtez le job pour éviter toute accumulation supplémentaire de fichiers :

yarn application -kill <appId>