O Realtime Compute for Apache Flink simplifica a ingestão de dados em tempo real ao gerenciar automaticamente a transição da sincronização completa para a incremental, a descoberta de metadados, a evolução de schema e a sincronização de todo o banco de dados. Este tópico mostra como criar rapidamente um job de ingestão que transmite dados do ApsaraDB RDS for MySQL para o Hologres.
Contexto
A figura a seguir ilustra como esses bancos de dados e tabelas aparecem no console DMS.
Siga estas etapas para desenvolver um job de ingestão que sincronize todas essas tabelas com o Hologres e consolide as tabelas de usuários fragmentadas em uma única tabela:
Este tópico utiliza a ingestão de dados Flink CDC para sincronizar todo o banco de dados e mesclar tabelas fragmentadas. Isso permite concluir a sincronização de dados completa e incremental, bem como a sincronização de alterações de schema em tempo real, com um único job.
Pré-requisitos
Se você utilizar um usuário RAM ou uma função RAM, verifique se possui as permissões necessárias para acessar o console Flink. Para obter mais informações, consulte Gerenciar permissões.
Crie um workspace Flink. Para obter mais informações, consulte Ativar o Realtime Compute for Apache Flink.
-
Armazenamento upstream e downstream
Crie uma instância ApsaraDB RDS for MySQL. Para obter mais informações, consulte (Descontinuado, redireciona para "Etapa 1") Criar rapidamente uma instância ApsaraDB RDS for MySQL.
Crie uma instância Hologres. Para obter mais informações, consulte Comprar uma instância Hologres.
NotaAs instâncias ApsaraDB RDS for MySQL e Hologres devem estar na mesma região e Virtual Private Cloud (VPC) do workspace Flink. Caso contrário, estabeleça uma conexão de rede. Para obter mais informações, consulte Como acesso outros serviços entre VPCs? e Como acesso a internet?.
Prepare dados de teste e configure listas de permissões de IP. Para obter mais informações, consulte Preparar dados de teste do MySQL e um banco de dados Hologres e Configurar uma lista de permissões de IP.
Preparar dados de teste do MySQL e um banco de dados Hologres
Clique em tpc_ds.sql, user_db1.sql, user_db2.sql e user_db3.sql para baixar os arquivos de dados de teste para sua máquina local.
-
No console DMS, prepare os dados de teste na instância ApsaraDB RDS for MySQL.
-
Faça login na instância ApsaraDB RDS for MySQL usando o DMS.
Para obter mais informações, consulte (Descontinuado, redireciona para "Etapa 2") Fazer login em uma instância ApsaraDB RDS for MySQL usando o DMS.
-
Na janela SQL Console, insira os comandos a seguir e clique em Execute.
Os comandos a seguir criam quatro bancos de dados:
tpc_ds,user_db1,user_db2euser_db3.CREATE DATABASE tpc_ds; CREATE DATABASE user_db1; CREATE DATABASE user_db2; CREATE DATABASE user_db3; Na barra de navegação superior, clique em Data Import.
-
Na aba Batch Data Import, selecione um banco de dados para importar os dados, faça upload do arquivo SQL correspondente, clique em Submit e, em seguida, clique em Execute Change. Na caixa de diálogo exibida, clique em Confirm Execution.
Repita essa operação para importar os arquivos de dados correspondentes aos bancos de dados
tpc_ds,user_db1,user_db2euser_db3.
-
-
No console Hologres, crie um banco de dados chamado
my_userpara armazenar os dados consolidados da tabela de usuários.Para obter mais informações, consulte Criar um banco de dados.
Configurar uma lista de permissões de IP
Para permitir que o Flink acesse as instâncias ApsaraDB RDS for MySQL e Hologres, adicione o bloco CIDR do seu workspace Flink à lista de permissões de IP de ambas as instâncias.
-
Obtenha o bloco CIDR do workspace Flink.
Faça login no console Realtime Compute for Apache Flink.
Na lista de workspaces, localize o workspace desejado e escolha na coluna Actions.
-
Na caixa de diálogo Workspace Details, visualize as informações de CIDR Block do vSwitch do Flink.

-
Adicione o bloco CIDR do Flink à lista de permissões de IP da instância ApsaraDB RDS for MySQL.
Para obter mais informações, consulte Configurar uma lista de permissões de IP.

-
Adicione o bloco CIDR do Flink à lista de permissões de IP da instância Hologres.
Ao configurar uma conexão de dados no HoloWeb, defina Login Method como Password-free login for current user antes de configurar uma lista de permissões de IP para a conexão. Para obter mais informações, consulte Lista de permissões de IP.

Etapa 1: Desenvolver o job de ingestão
-
Faça login no console de desenvolvimento Flink e crie um novo job.
Na página , clique em New.
-
Clique em Blank Data Ingestion Draft.
O Realtime Compute for Apache Flink oferece um conjunto abrangente de modelos de código, cada um com casos de uso específicos, exemplos de código e orientações. Clique em um modelo para conhecer os recursos do produto e a sintaxe necessária para implementar sua lógica de negócios.
Clique em Next.
-
Na caixa de diálogo New Data Ingestion Job Draft, configure os parâmetros do job.
Parâmetro
Descrição
Exemplo
File Name
Nome do job.
NotaO nome do job deve ser exclusivo no projeto atual.
flink-test
Storage Location
Pasta onde o arquivo de código do job será armazenado.
Também é possível clicar no ícone
ao lado de uma pasta existente para criar uma subpasta.Job Drafts
Engine Version
Versão do mecanismo Flink utilizada pelo job. Para obter informações sobre números de versão do mecanismo, compatibilidade de versões e datas importantes do ciclo de vida, consulte Versões do mecanismo.
vvr-11,1-jdk11-flink-1,20
Clique em OK.
-
Copie o código do job a seguir para o editor de jobs.
O código abaixo sincroniza todas as tabelas do banco de dados
tpc_dscom o Hologres e consolida as tabelas de usuários fragmentadas em uma única tabela no Hologres:source: type: mysql name: MySQL Source hostname: localhost port: 3306 username: username password: password tables: tpc_ds.\.*,user_db[0-9]+.user[0-9]+ server-id: 8601-8604 # (Optional) Synchronize table and column comments. include-comments.enabled: true # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory errors. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to accelerate reading. scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: hologres name: Hologres Sink endpoint: ****.hologres.aliyuncs.com:80 dbname: cdcyaml_test username: ${secret_values.holo-username} password: ${secret_values.holo-password} sink.type-normalize-strategy: BROADEN route: # Merge and synchronize the sharded user tables to the my_user.users table. - source-table: user_db[0-9]+.user[0-9]+ sink-table: my_user.usersNotaAs tabelas do banco de dados MySQL
tpc_dssão mapeadas diretamente para tabelas com nomes idênticos no destino, portanto, a seçãoroutenão exige configuração extra de mapeamento. Para sincronizar as tabelas com um banco de dados de nome diferente, comoods_tps_ds, configure o módulorouteda seguinte forma:route: # Merge and synchronize the sharded user tables to the my_user.users table. - source-table: user_db[0-9]+.user[0-9]+ sink-table: my_user.users # Rename the database for all tables under tpc_ds and synchronize them to ods_tps_ds. - source-table: tpc_ds.\.* sink-table: ods_tps_ds.<> replace-symbol: <>
Etapa 2: Iniciar o job
Na página , clique em Deploy. Na caixa de diálogo exibida, clique em Confirm.

Na página , clique em Start na coluna Actions do job desejado. Configure os parâmetros conforme necessário. Para obter mais informações, consulte Iniciar um job.
-
Clique em Start.
Após o início do job, visualize suas informações de execução e status na página Job Operations.

Etapa 3: Verificar a sincronização completa
Faça login no Console de Gerenciamento do Hologres.
-
Na aba Metadata Management, verifique se as 24 tabelas e seus dados estão presentes no banco de dados
tpc_dsda instância Hologres.
-
Na aba Metadata Management, verifique o schema da tabela
usersno banco de dadosmy_user.As figuras a seguir mostram o schema e os dados da tabela sincronizada.
-
Schema da tabela

O schema da tabela
usersinclui duas colunas adicionais não encontradas nas tabelas MySQL de origem:_db_namee_table_name. Essas colunas indicam o banco de dados e a tabela de origem de cada linha e fazem parte da chave primária composta, garantindo a unicidade dos dados após a consolidação das tabelas fragmentadas. -
Dados da tabela
No canto superior direito da página de informações da tabela
users, clique em Query Table. Insira o comando a seguir e clique em Run.select * from users order by _db_name,_table_name,id;O resultado da consulta é mostrado na figura a seguir.

-
Etapa 4: Verificar a sincronização incremental
Após a conclusão da sincronização completa, o job muda automaticamente para a fase de sincronização incremental sem intervenção manual. Verifique o valor de currentEmitEventTimeLag na aba Monitoring and Alerts para determinar a fase de sincronização de dados.
Faça login no console Realtime Compute for Apache Flink.
Clique em Console na coluna Actions do workspace desejado.
Na página , clique no nome do job desejado.
Clique na aba Monitoring and Alerts (ou Metrics).
-
Analise o gráfico
currentEmitEventTimeLagpara determinar a fase de sincronização de dados.
Um valor igual a 0 indica a fase de sincronização completa.
Um valor maior que 0 indica a fase de sincronização incremental.
-
Verifique a sincronização de dados e alterações de schema em tempo real.
A source MySQL CDC suporta sincronização de dados e schema em tempo real durante a fase incremental. Para verificar isso, modifique o schema e os dados de uma tabela de usuários fragmentada no MySQL após o job entrar nessa fase.
-
Faça login na instância ApsaraDB RDS for MySQL usando o DMS.
Para obter mais informações, consulte (Descontinuado, redireciona para "Etapa 2") Fazer login em uma instância ApsaraDB RDS for MySQL usando o DMS.
-
No banco de dados
user_db2, execute os comandos a seguir para modificar o schema da tabelauser02e para inserir e atualizar dados.USE `user_db2`; ALTER TABLE `user02` ADD COLUMN `age` INT; -- Add the age column. INSERT INTO `user02` (id, name, age) VALUES (27, 'Tony', 30); -- Insert a row that includes the age data. UPDATE `user05` SET name='JARK' WHERE id=15; -- Update another table and change the name to uppercase. -
No console Hologres, verifique as alterações no schema e nos dados da tabela
users.No canto superior direito da página de informações da tabela
users, clique em Query Table, insira o comando a seguir e clique em Run.select * from users order by _db_name,_table_name,id;A figura a seguir mostra o resultado da consulta. A alteração de schema em
user02e as modificações de dados são propagadas em tempo real, mesmo que as tabelas fragmentadas tenham schemas diferentes. A tabelausersdo Hologres agora exibe a nova colunaage, o registro inserido para Tony e o registro atualizado para JARK.
-
(Opcional) Etapa 5: Configurar recursos do job
Para obter melhor desempenho, ajuste os recursos do job, como concorrência, memória do TaskManager e CUs, com base no volume de dados.
Na página , clique no nome do job desejado.
Na aba Deployment Details, clique em Edit no canto superior direito da seção Resource Configuration.
Defina manualmente parâmetros de recursos, como memória do TaskManager e concorrência.
No lado direito da seção Resource Configuration, clique em Save.
-
Reinicie o job.
As alterações na configuração de recursos só entram em vigor após o reinício do job.
Documentos relacionados
Para a sintaxe de cada módulo de ingestão de dados, consulte Referência de desenvolvimento de job de ingestão de dados Flink CDC.
Se encontrar problemas durante a execução do job de ingestão de dados, consulte Problemas comuns e soluções para jobs de ingestão de dados.