Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Ingestão de banco de dados em tempo real

Última atualização: Jun 27, 2026

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

Preparar dados de teste do MySQL e um banco de dados Hologres

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

  2. No console DMS, prepare os dados de teste na instância ApsaraDB RDS for MySQL.

    1. Faça login na instância ApsaraDB RDS for MySQL usando o DMS.

    2. 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_db2 e user_db3.

      CREATE DATABASE tpc_ds;
      CREATE DATABASE user_db1;
      CREATE DATABASE user_db2;
      CREATE DATABASE user_db3;
    3. Na barra de navegação superior, clique em Data Import.

    4. 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_db2 e user_db3.导入数据

  3. No console Hologres, crie um banco de dados chamado my_user para 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.

  1. Obtenha o bloco CIDR do workspace Flink.

    1. Faça login no console Realtime Compute for Apache Flink.

    2. Na lista de workspaces, localize o workspace desejado e escolha More > Workspace Details na coluna Actions.

    3. Na caixa de diálogo Workspace Details, visualize as informações de CIDR Block do vSwitch do Flink.

      网段信息

  2. 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.RDS白名单

  3. 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.Holo白名单

Etapa 1: Desenvolver o job de ingestão

  1. Faça login no console de desenvolvimento Flink e crie um novo job.

    1. Na página Data Development > Data Ingestion, clique em New.

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

    3. Clique em Next.

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

      Nota

      O 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

    5. Clique em OK.

  2. 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_ds com 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.users
    Nota

    As tabelas do banco de dados MySQL tpc_ds são mapeadas diretamente para tabelas com nomes idênticos no destino, portanto, a seção route não exige configuração extra de mapeamento. Para sincronizar as tabelas com um banco de dados de nome diferente, como ods_tps_ds, configure o módulo route da 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

  1. Na página Data Development > Data Ingestion, clique em Deploy. Na caixa de diálogo exibida, clique em Confirm.部署

  2. Na página Operation Center > Job Operations, 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.

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

  1. Faça login no Console de Gerenciamento do Hologres.

  2. Na aba Metadata Management, verifique se as 24 tabelas e seus dados estão presentes no banco de dados tpc_ds da instância Hologres.

    holo表数据

  3. Na aba Metadata Management, verifique o schema da tabela users no banco de dados my_user.

    As figuras a seguir mostram o schema e os dados da tabela sincronizada.

    • Schema da tabela 表结构

      O schema da tabela users inclui duas colunas adicionais não encontradas nas tabelas MySQL de origem: _db_name e _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.

  1. Faça login no console Realtime Compute for Apache Flink.

  2. Clique em Console na coluna Actions do workspace desejado.

  3. Na página Operation Center > Job Operations, clique no nome do job desejado.

  4. Clique na aba Monitoring and Alerts (ou Metrics).

  5. Analise o gráfico currentEmitEventTimeLag para 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.

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

    1. Faça login na instância ApsaraDB RDS for MySQL usando o DMS.

    2. No banco de dados user_db2, execute os comandos a seguir para modificar o schema da tabela user02 e 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.
    3. 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 user02 e as modificações de dados são propagadas em tempo real, mesmo que as tabelas fragmentadas tenham schemas diferentes. A tabela users do Hologres agora exibe a nova coluna age, 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.

  1. Na página Operation Center > Job Operations, clique no nome do job desejado.

  2. Na aba Deployment Details, clique em Edit no canto superior direito da seção Resource Configuration.

  3. Defina manualmente parâmetros de recursos, como memória do TaskManager e concorrência.

  4. No lado direito da seção Resource Configuration, clique em Save.

  5. Reinicie o job.

    As alterações na configuração de recursos só entram em vigor após o reinício do job.

Documentos relacionados