O Apache Flink open source não oferece suporte a gravações em streaming no OSS-HDFS nem à semântica exactly-once para mídias de armazenamento. Para ative esses recursos, use o JindoSDK com o Apache Flink.
Para gravar dados em streaming no OSS-HDFS sem implantar o JindoSDK, use o Realtime Compute for Apache Flink. Para mais informações, consulte 实时计算Flink读写OSS或者OSS-HDFS.
Pré-requisitos
Crie uma instância ECS. Create an instance.
Ative o OSS-HDFS para um bucket e conceda as permissões de acesso adequadas. Para mais informações, consulte 开通OSS-HDFS服务.
Instale o Apache Flink 1.10.1 ou posterior. As versões 1.16.0 e posteriores do Flink não foram verificadas. Baixe o Flink em Apache Flink.
Configure o JindoSDK
Conecte-se à instância ECS. Connect to an instance.
Baixe e descompacte a versão mais recente do pacote JAR do JindoSDK. Para baixe o JindoSDK, visite GitHub.
-
Mova o arquivo jindo-flink-${version}-full.jar do diretório plugins/flink/ do JindoSDK para o diretório lib raiz do Apache Flink.
mv plugins/flink/jindo-flink-${version}-full.jar lib/
Se o Apache Flink tiver um conector OSS integrado, remova-o excluindo o arquivo
flink-oss-fs-hadoop-${flink-version}.jardo subdiretóriolibou do caminhoplugins/oss-fs-hadoop.Após configure o JindoSDK, use o prefixo
oss://nos jobs do Flink para gravar dados em streaming no OSS-HDFS ou no OSS. O JindoSDK identifica automaticamente o armazenamento de destino.
Exemplos
-
Configure as definições gerais.
Para gravar dados no OSS-HDFS com semântica exactly-once:
-
Ative o checkpointing do Apache Flink.
Código de exemplo:
-
Crie um StreamExecutionEnvironment:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); -
Ative o checkpointing:
env.enableCheckpointing(<userDefinedCheckpointInterval>, CheckpointingMode.EXACTLY_ONCE);
-
Use uma fonte de dados que ofereça suporte a retransmissão, como o Kafka.
-
-
Defina as configurações rápidas.
Especifique um caminho iniciado por oss:// com buckets e endpoints do OSS-HDFS. Nenhuma dependência adicional é necessária.
-
Adicione um sink.
O código a seguir grava um DataStream<String> no 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);ImportanteO campo
.{oss-hdfs-endpoint}é opcional. Se omitido, especifique o endpoint correto do OSS-HDFS na configuração do Flink ou do Hadoop. Execute
env.execute()para envie o job do Flink.
-
(Opcional) Defina configurações personalizadas
Defina parâmetros personalizados ao envie jobs do Flink para ative recursos específicos.
Use -yD para configure envios de jobs do Flink baseados em YARN:
<flink_home>/bin/flink run -m yarn-cluster -yD key1=value1 -yD key2=value2 ...
A injeção de entropia substitui uma string especificada no caminho de destino por uma string aleatória para distribuir as gravações entre partições e melhorar o desempenho.
Configure os seguintes parâmetros para gravações no OSS-HDFS:
oss.entropy.key=<user-defined-key>
oss.entropy.length=<user-defined-length>
Durante as gravações, a string <user-defined-key> no caminho é substituída por uma string aleatória cujo comprimento equivale a <user-defined-length>. O valor de <user-defined-length> deve ser maior que 0.