Todos os produtos
Search
Central de documentação

Object Storage Service:Use o Apache Flink para gravar dados no OSS-HDFS

Última atualização: Jul 03, 2026

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.

Nota

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

  1. Conecte-se à instância ECS. Connect to an instance.

  2. Baixe e descompacte a versão mais recente do pacote JAR do JindoSDK. Para baixe o JindoSDK, visite GitHub.

  3. 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/
Importante
  • Se o Apache Flink tiver um conector OSS integrado, remova-o excluindo o arquivo flink-oss-fs-hadoop-${flink-version}.jar do subdiretório lib ou do caminho plugins/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

  1. Configure as definições gerais.

    Para gravar dados no OSS-HDFS com semântica exactly-once:

    1. Ative o checkpointing do Apache Flink.

      Código de exemplo:

      1. Crie um StreamExecutionEnvironment:

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
      2. Ative o checkpointing:

        env.enableCheckpointing(<userDefinedCheckpointInterval>, CheckpointingMode.EXACTLY_ONCE);
    2. Use uma fonte de dados que ofereça suporte a retransmissão, como o Kafka.

  2. 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.

    1. 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);
      Importante

      O campo .{oss-hdfs-endpoint} é opcional. Se omitido, especifique o endpoint correto do OSS-HDFS na configuração do Flink ou do Hadoop.

    2. 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.