Tous les produits
Search
Centre de documentation

Object Storage Service:Use Flink on an EMR cluster to write data to the OSS-HDFS service

Dernière mise à jour :Aug 18, 2026

La fonctionnalité d'écriture avec reprise permet d'écrire des données sur les supports de stockage en utilisant la sémantique EXACTLY_ONCE. Cette rubrique décrit comment utiliser Apache Flink sur un cluster E-MapReduce (EMR) pour écrire des données dans OSS-HDFS avec prise en charge de la reprise.

Exemple

  1. Configurations générales

    Pour écrire des données dans le service OSS-HDFS en utilisant la sémantique EXACTLY_ONCE, vous devez effectuer les configurations suivantes :

    1. Activez les points de contrôle (checkpoints) pour Flink.

      Exemple :

      1. Créez un objet StreamExecutionEnvironment comme suit.

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
      2. Exécutez la commande suivante pour activer les points de contrôle.

        env.enableCheckpointing(<userDefinedCheckpointInterval>, CheckpointingMode.EXACTLY_ONCE);
    2. Utilisez une source de données pouvant être rejouée, telle que Kafka.

  2. Utilisation simple

    Aucune dépendance supplémentaire n'est requise. Pour écrire des données, utilisez un chemin commençant par le préfixe oss:// et spécifiez le bucket et l'Endpoint du service OSS-HDFS.

    1. Ajoutez un sink.

      L'exemple suivant montre comment écrire un objet DataStream<String> dans le service OSS-HDFS.

      String outputPath = "oss://{user-defined-oss-hdfs-bucket.oss-hdfs-endpoint}/{user-defined-dir}";
      StreamingFileSink<String> sink = StreamingFileSink.forRowFormat(
              new Path(outputPath),
              new SimpleStringEncoder<String>("UTF-8")
      ).build();
      outputStream.addSink(sink);
      Important

      La partie .{oss-hdfs-endpoint} du chemin est facultative. Si vous omettez cette partie, vous devez configurer correctement l'Endpoint du service OSS-HDFS dans le composant Flink ou Hadoop.

    2. Utilisez env.execute() pour exécuter le job Flink.

(Facultatif) Configurer des paramètres personnalisés

Vous pouvez configurer des paramètres personnalisés lors de la soumission des jobs Flink pour activer des fonctionnalités spécifiques.

Utilisez -yD pour configurer les soumissions de jobs Flink basées sur YARN :

<flink_home>/bin/flink run -m yarn-cluster -yD key1=value1 -yD key2=value2 ...

L'injection d'entropie remplace une chaîne spécifiée dans le chemin de destination par une chaîne aléatoire afin de distribuer les écritures sur plusieurs partitions et d'améliorer les performances.

Configurez les paramètres suivants pour les écritures dans OSS-HDFS :

oss.entropy.key=<user-defined-key>
oss.entropy.length=<user-defined-length>

Lors des écritures, la chaîne <user-defined-key> dans le chemin est remplacée par une chaîne aléatoire dont la longueur est égale à <user-defined-length>. La valeur de <user-defined-length> doit être supérieure à 0.