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
-
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 :
-
Activez les points de contrôle (checkpoints) pour Flink.
Exemple :
-
Créez un objet StreamExecutionEnvironment comme suit.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); -
Exécutez la commande suivante pour activer les points de contrôle.
env.enableCheckpointing(<userDefinedCheckpointInterval>, CheckpointingMode.EXACTLY_ONCE);
-
Utilisez une source de données pouvant être rejouée, telle que Kafka.
-
-
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.
-
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);ImportantLa 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. 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.