Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Flink CDC data ingestion job

Última atualização: Jun 27, 2026

O Realtime Compute for Apache Flink oferece um recurso avançado de ingestão de dados baseado no Flink CDC. Este guia demonstra como crie um job de ingestão de dados do Flink CDC para sincronizar um banco de dados MySQL inteiro com um banco de dados StarRocks.

Pré-requisitos

Informações de fundo

Considere que sua instância do ApsaraDB RDS for MySQL possua um banco de dados chamado order_dw_mysql contendo três tabelas de negócios: orders, orders_pay e product_catalog. Para sincronizar essas tabelas e seus dados com o banco de dados order_dw_sr no StarRocks, siga estas etapas:

  1. Etapa 1: Preparar dados de teste no ApsaraDB RDS for MySQL

  2. Etapa 2: Desenvolver um job de ingestão de dados do Flink CDC

  3. Etapa 3: Iniciar o job de ingestão de dados do Flink CDC

  4. Etapa 4: Verificar os resultados da sincronização no StarRocks

Etapa 1: Preparar dados de teste do MySQL

  1. Crie um banco de dados e uma conta.

    Crie um banco de dados chamado order_dw_mysql e uma conta padrão com permissões de leitura e gravação nele. Para mais informações, consulte (Descontinuado, redirecionado para "Etapa 1") Criar um banco de dados e uma conta e Gerenciar bancos de dados.

  2. Faça login na instância do ApsaraDB RDS for MySQL usando o Data Management (DMS).

  3. Na janela do SQL Console, insira os comandos abaixo e clique em Execute para criar três tabelas de negócios e inserir dados.

    CREATE TABLE `orders` (
      order_id bigint not null primary key,
      user_id varchar(50) not null,
      shop_id bigint not null,
      product_id bigint not null,
      buy_fee numeric(20,2) not null,   
      create_time timestamp not null,
      update_time timestamp not null default now(),
      state int not null 
    );
    CREATE TABLE `orders_pay` (
      pay_id bigint not null primary key,
      order_id bigint not null,
      pay_platform int not null, 
      create_time timestamp not null
    );
    CREATE TABLE `product_catalog` (
      product_id bigint not null primary key,
      catalog_name varchar(50) not null
    );
    -- Prepare data
    INSERT INTO product_catalog VALUES(1, 'phone_aaa'),(2, 'phone_bbb'),(3, 'phone_ccc'),(4, 'phone_ddd'),(5, 'phone_eee');
    INSERT INTO orders VALUES
    (100001, 'user_001', 12345, 1, 5000.05, '2023-02-15 16:40:56', '2023-02-15 18:42:56', 1),
    (100002, 'user_002', 12346, 2, 4000.04, '2023-02-15 15:40:56', '2023-02-15 18:42:56', 1),
    (100003, 'user_003', 12347, 3, 3000.03, '2023-02-15 14:40:56', '2023-02-15 18:42:56', 1),
    (100004, 'user_001', 12347, 4, 2000.02, '2023-02-15 13:40:56', '2023-02-15 18:42:56', 1),
    (100005, 'user_002', 12348, 5, 1000.01, '2023-02-15 12:40:56', '2023-02-15 18:42:56', 1),
    (100006, 'user_001', 12348, 1, 1000.01, '2023-02-15 11:40:56', '2023-02-15 18:42:56', 1),
    (100007, 'user_003', 12347, 4, 2000.02, '2023-02-15 10:40:56', '2023-02-15 18:42:56', 1);
    INSERT INTO orders_pay VALUES
    (2001, 100001, 1, '2023-02-15 17:40:56'),
    (2002, 100002, 1, '2023-02-15 17:40:56'),
    (2003, 100003, 0, '2023-02-15 17:40:56'),
    (2004, 100004, 0, '2023-02-15 17:40:56'),
    (2005, 100005, 0, '2023-02-15 18:40:56'),
    (2006, 100006, 0, '2023-02-15 18:40:56'),
    (2007, 100007, 0, '2023-02-15 18:40:56');

Etapa 2: Desenvolver um job do Flink CDC

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

  2. Clique em Console para acessar o workspace do projeto.

  3. No painel de navegação à esquerda, escolha Development > Data Ingestion.

  4. Clique no ícone image, clique em New Draft with Template, selecione MySQL to StarRocks data synchronization e clique em Next.

  5. Insira um Job Name e um Storage Location, selecione uma Engine Version e clique em OK.

  6. Configure o código YAML do job.

    O código a seguir apresenta um exemplo de sincronização de todas as tabelas do banco de dados order_dw_mysql no MySQL para o banco de dados order_dw_sr no StarRocks.

    source:
      type: mysql
      hostname: rm-bp1rk934iidc3****.mysql.rds.aliyuncs.com
      port: 3306
      username: ${secret_values.mysqlusername}
      password: ${secret_values.mysqlpassword}
      tables: order_dw_mysql.\.*
      server-id: 8601-8604
      # (Optional) Synchronize data from tables that are newly created during the incremental phase.
      scan.binlog.newly-added-table.enabled: true
      # (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: starrocks
      name: StarRocks Sink
      jdbc-url: jdbc:mysql://fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:9030
      load-url: fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:8030
      username: ${secret_values.starrocksusername}
      password: ${secret_values.starrockspassword}
      table.create.properties.replication_num: 1
      sink.buffer-flush.interval-ms: 5000 # Flush data every 5 seconds.
    route:
      - source-table: order_dw_mysql.\.*
        sink-table: order_dw_sr.<>
        replace-symbol: <>
        description: route all tables in source_db to sink_db
    pipeline:
      name: Sync MySQL Database to StarRocks

    A tabela a seguir descreve os parâmetros de configuração necessários para este exemplo. Para mais detalhes sobre os parâmetros de ingestão de dados, consulte MySQL e StarRocks.

    Nota

    Jobs YAML suportam apenas variáveis de projeto. Utilize variáveis para evitar que informações como senhas sejam exibidas em texto simples. Para mais informações, consulte Gerenciamento de Variáveis.

    Categoria

    Parâmetro

    Descrição

    Valor de exemplo

    source

    hostname

    Endereço IP ou nome do host do banco de dados MySQL.

    Recomendamos o uso do endpoint interno.

    rm-bp1rk934iidc3****.mysql.rds.aliyuncs.com

    port

    Número da porta do serviço de banco de dados MySQL.

    3306

    username

    Nome de usuário e senha do banco de dados MySQL. Utilize as credenciais da conta criada na Etapa 1: Preparar dados de teste no ApsaraDB RDS for MySQL.

    ${secret_values.mysqlusername}

    password

    ${secret_values.mysqlpassword}

    tables

    Nomes das tabelas do MySQL. É possível usar expressões regulares para ler dados de várias tabelas.

    Neste tópico, todas as tabelas e dados do banco de dados order_dw_mysql são sincronizados.

    order_dw_mysql.\.*

    server-id

    ID numérico exclusivo para a conexão do cliente ao banco de dados.

    8601-8604

    sink

    jdbc-url

    URL de conexão JDBC.

    Especifique o endereço IP e a porta de consulta do Frontend (FE) no formato jdbc:mysql://ip:port.

    Na aba Instance Details do console E-MapReduce, visualize o internal endpoint e a query port do FE da instância de destino.

    jdbc:mysql://fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:9030

    load-url

    URL do serviço HTTP usada para conectar ao nó FE.

    Na aba Instance Details do console E-MapReduce, visualize o internal endpoint e a HTTP port do FE da instância de destino.

    fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:8030

    username

    Credenciais para conexão com o StarRocks.

    Utilize as credenciais configuradas durante a criação da instância do StarRocks.

    Nota

    Este exemplo utiliza variáveis para evitar a exposição de credenciais em texto simples. Para mais informações, consulte Gerenciar variáveis.

    ${secret_values.starrocksusername}

    password

    ${secret_values.starrockspassword}

    sink.buffer-flush.interval-ms

    Intervalo de liberação (flush) do buffer interno.

    Um intervalo curto (5 segundos) é utilizado porque este exemplo contém poucos dados, permitindo visualizar os resultados rapidamente.

    5000

    route

    source-table

    Tabela ou tabelas de origem a serem roteadas.

    Use uma expressão regular para corresponder a várias tabelas. Por exemplo, order_dw_mysql.\.* roteia todas as tabelas do banco de dados order_dw_mysql.

    order_dw_mysql.\.*

    sink-table

    Padrão da tabela de destino para os dados roteados.

    Utilize o símbolo definido no parâmetro replace-symbol como espaço reservado para cada nome de tabela de origem, possibilitando o roteamento muitos-para-muitos.

    Para mais informações sobre regras de roteamento, consulte Módulo Route.

    order_dw_sr.<>

    replace-symbol

    Espaço reservado para o nome da tabela de origem usado na correspondência de padrões.

    <>

  7. Clique em Deploy.

Etapa 3: Iniciar o job do Flink CDC

  1. Na página Data Ingestion, clique em Deploy e, em seguida, clique em OK na caixa de diálogo exibida.

  2. Na página O&M > Deployments, localize o job YAML desejado e clique em Start na coluna Actions.

  3. Clique em Start.

    Neste exemplo, selecione Initial Mode. Para mais informações sobre os parâmetros, consulte Iniciar um job. Após o início do job, monitore seu status na página Deployments.

Etapa 4: Verificar resultados no StarRocks

Depois que o job entrar no estado RUNNING, verifique os dados no StarRocks.

  1. Conectar a uma instância do EMR Serverless StarRocks usando o EMR StarRocks Manager.

  2. No painel de navegação à esquerda, clique em SQL Editor. Na aba Database, clique no ícone de atualização image.

    Um banco de dados chamado order_dw_sr aparecerá sob default_catalog.

  3. Na aba Query List, clique em + File para criar um Query Script. Insira as instruções SQL a seguir e clique em Run.

    SELECT * FROM default_catalog.order_dw_sr.orders order by order_id;
    SELECT * FROM default_catalog.order_dw_sr.orders_pay order by pay_id;
    SELECT * FROM default_catalog.order_dw_sr.product_catalog order by product_id;
  4. Visualize os resultados abaixo dos comandos.

    Os resultados mostram que as tabelas e os dados do banco de dados MySQL agora existem no StarRocks.

    As tabelas sincronizadas incluem default_catalog.order_dw_sr.orders, default_catalog.order_dw_sr.orders_pay e default_catalog.order_dw_sr.product_catalog. Execute instruções SELECT para consultar cada tabela e verificar a integridade dos dados.

Documentação relacionada