Este tópico descreve como desenvolver um job do Realtime Compute for Apache Flink que usa a DataStream API para gravar dados em um catálogo do Data Lake Formation (DLF) via Paimon REST.
Pré-requisitos
Você tem um workspace totalmente gerenciado do Realtime Compute for Apache Flink. Caso ainda não tenha criado um, consulte Activate Realtime Compute for Apache Flink.
Um catálogo DLF foi criado. Para obter mais informações, consulte Get started with DLF.
Seu workspace do Flink e o catálogo DLF estão na mesma região.
-
A VPC do seu workspace do Flink está na lista de permissões de VPC do DLF. Para obter mais informações, consulte API usage guide.
NotaO DLF ativa o acesso por VPC por padrão. Para ativar o acesso por rede pública, consulte Ative and configure public network access.
Preparações
Baixe o JAR empacotado do Paimon
paimon-flink-*.jarna versão 1.1 ou posterior no site do Apache Paimon.Baixe o
paimon-oss-*.jarna versão 1.1 ou posterior em Apache Paimon Filesystems.
Escolha um método de dependência
O ambiente de runtime do Flink não inclui o conector Paimon nem o sistema de arquivos OSS. Use um dos métodos a seguir para garantir que paimon-flink-*.jar e paimon-oss-*.jar estejam disponíveis durante a execução do job.
Method 1: Upload additional files in the console
Não é necessário modifique o arquivo pom.xml. Ao crie um job JAR no console de desenvolvimento do Realtime Compute for Apache Flink, faça o upload dos arquivos paimon-flink-*.jar e paimon-oss-*.jar baixados na seção Preparations como arquivos de dependência adicionais.
Method 2: Package dependencies into a fat JAR with Maven
Adicione as seguintes dependências e propriedades ao arquivo pom.xml do seu projeto.
<properties>
<!-- Paimon version. Specify 1.1 or later. -->
<paimon.version>1.1.0</paimon.version>
<!-- Flink major version. Set based on your VVR version. See the following table. -->
<flink.main.version>1.20</flink.main.version>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.paimon</groupId>
<artifactId>paimon-flink-${flink.main.version}</artifactId>
<version>${paimon.version}</version>
</dependency>
<dependency>
<groupId>org.apache.paimon</groupId>
<artifactId>paimon-oss</artifactId>
<version>${paimon.version}</version>
</dependency>
</dependencies>
O valor de ${flink.main.version} é defina conforme a tabela a seguir.
|
Versão do VVR |
|
|
VVR 8.x |
1,17 |
|
VVR 11.x |
1,20 |
Com este método, as dependências são empacotadas no fat JAR, dispensando o upload de arquivos JAR adicionais durante a implantação.
Etapa 1: Escreva o código do job
No método main() do seu job DataStream, use o código a seguir para crie uma instância de catálogo DLF.
Options options = new Options();
options.set("type", "paimon");
options.set("metastore", "rest");
options.set("uri", "http://<region-id>-vpc.dlf.aliyuncs.com");
options.set("warehouse", "your-catalog-name");
options.set("token.provider", "dlf");
options.set("dlf.access-key-id", "your-access-key-id");
options.set("dlf.access-key-secret", "your-access-key-secret");
Catalog catalog = FlinkCatalogFactory.createPaimonCatalog(options);
Parâmetros obrigatórios:
|
Parâmetro |
Descrição |
Exemplo |
|
|
O tipo de catálogo, analisado automaticamente a partir do JAR personalizado. Não altere este valor. |
|
|
|
O tipo de metastore para o DLF. Defina como |
|
|
|
O endpoint VPC do servidor de catálogo REST do DLF. O formato é |
|
|
|
O nome do catálogo Paimon. |
|
|
|
O provedor de token. Defina como |
|
|
|
O ID do seu AccessKey para autenticação. Para obter mais informações, consulte Visualize RAM user AccessKey information. |
|
|
|
O segredo do seu AccessKey para autenticação. |
Após criar o catálogo, registre-o e utilize-o no seu job DataStream para ler e gravar tabelas Paimon.
Etapa 2: Empacote e implante o job
Empacote seu job DataStream em um arquivo JAR.
Faça o upload do JAR do job no console do Realtime Compute for Apache Flink e envie o job.
Se você optou pelo Method 1 (upload de arquivos adicionais no console), adicione
paimon-flink-*.jarepaimon-oss-*.jaràs dependências adicionais ao envie o job.
Para obter mais informações sobre como desenvolver e depurar jobs Flink JAR, consulte Develop a JAR job.