Tous les produits
Search
Centre de documentation

Object Storage Service:Utiliser Apache Flink pour écrire des données dans OSS-HDFS

Dernière mise à jour :Aug 18, 2026

La version open source d'Apache Flink ne prend pas en charge l'écriture en streaming vers OSS-HDFS ni la sémantique exactly-once pour les supports de stockage. Pour activer ces fonctionnalités, utilisez JindoSDK avec Apache Flink.

Remarque

Pour écrire des données en streaming dans OSS-HDFS sans déployer JindoSDK, utilisez Realtime Compute for Apache Flink. Pour plus d'informations, consultez la rubrique Lecture et écriture de données dans OSS ou OSS-HDFS avec Realtime Compute for Apache Flink.

Prérequis

  • Créez une instance ECS. Consultez la section Créer une instance.

  • Activez le service OSS-HDFS pour un bucket et accordez les autorisations d'accès appropriées. Pour plus d'informations, consultez la tâche Activer le service OSS-HDFS.

  • Installez Apache Flink 1.10.1 ou une version ultérieure. Les versions 1.16.0 et supérieures n'ont pas été vérifiées. Téléchargez Flink depuis la page Apache Flink.

Configurer JindoSDK

  1. Connectez-vous à l'instance ECS. Consultez la section Se connecter à une instance.

  2. Téléchargez et décompressez la dernière version du package JAR JindoSDK. Pour télécharger JindoSDK, accédez à GitHub.

  3. Déplacez le fichier jindo-flink-${version}-full.jar depuis le répertoire plugins/flink/ de JindoSDK vers le répertoire lib racine d'Apache Flink.

    mv plugins/flink/jindo-flink-${version}-full.jar lib/
Important
  • Si Apache Flink dispose d'un connecteur OSS intégré, supprimez-le en retirant le fichier flink-oss-fs-hadoop-${flink-version}.jar du sous-répertoire lib ou du chemin plugins/oss-fs-hadoop.

  • Après avoir configuré JindoSDK, utilisez le préfixe oss:// dans les jobs Flink pour écrire des données en streaming dans OSS-HDFS ou OSS. JindoSDK identifie automatiquement le stockage cible.

Exemples

  1. Configurez les paramètres généraux.

    Pour écrire des données dans OSS-HDFS avec une sémantique exactly-once :

    1. Activez la fonctionnalité de checkpointing d'Apache Flink.

      Exemple de code :

      1. Créez un StreamExecutionEnvironment :

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
      2. Activez le checkpointing :

        env.enableCheckpointing(<userDefinedCheckpointInterval>, CheckpointingMode.EXACTLY_ONCE);
    2. Utilisez une source de données qui prend en charge la retransmission des données, telle que Kafka.

  2. Configurez les paramètres rapides.

    Spécifiez un chemin commençant par oss:// avec les buckets et les endpoints OSS-HDFS. Aucune dépendance supplémentaire n'est requise.

    1. Ajoutez un sink.

      Le code suivant écrit un DataStream<String> dans 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

      Le champ .{oss-hdfs-endpoint} est facultatif. S'il est omis, spécifiez l'endpoint OSS-HDFS correct dans la configuration Flink ou Hadoop.

    2. Exécutez env.execute() pour soumettre 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 répartir les écritures sur les 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>

Pendant les é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.