Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Desenvolver um job de ingestão de dados Flink CDC

Última atualização: Jun 27, 2026

O Realtime Compute for Apache Flink usa o Flink CDC para ingestão de dados. Desenvolva jobs YAML para sincronizar dados eficientemente de uma source para um sink. Este tópico descreve como desenvolver um job de ingestão de dados Flink CDC.

Informações básicas

A ingestão de dados com Flink CDC usa o Flink CDC para simplificar a integração de dados. Ao usar YAML para definir processos ETL complexos, convertidos automaticamente em lógica de execução do Flink, você implementa recursos como sincronização completa de banco de dados, sincronização de tabelas fragmentadas, evolução de schema e colunas calculadas. Essa abordagem simplifica significativamente o processo de integração de dados, além de melhorar sua eficiência e confiabilidade.

Vantagens do Flink CDC

No Realtime Compute for Apache Flink, desenvolva um job de ingestão de dados Flink CDC, um job SQL ou um job DataStream para sincronizar dados. As seções a seguir descrevem as vantagens dos jobs de ingestão de dados Flink CDC em relação aos outros dois métodos de desenvolvimento.

Flink CDC vs. Flink SQL

Os jobs de ingestão de dados Flink CDC e os jobs Flink SQL transmitem tipos diferentes de dados:

  • Jobs SQL transmitem RowData, em que cada linha possui um tipo de alteração: insert (+I), update before (-U), update after (+U) ou delete (-D).

  • Jobs Flink CDC usam SchemaChangeEvent para transmitir informações de evolução de schema, como criação de tabela, adição de coluna ou truncamento de tabela. Eles utilizam DataChangeEvent para transmitir alterações de dados, como inserções, atualizações e exclusões. Uma mensagem de atualização contém os estados anterior e posterior de uma linha, o que permite gravar os dados de alteração originais no sink.

A tabela a seguir lista as vantagens dos jobs de ingestão de dados Flink CDC sobre os jobs SQL.

Flink CDC

Flink SQL

Descobre schemas automaticamente e suporta sincronização completa de banco de dados

Requer instruções manuais CREATE TABLE e INSERT

Suporta múltiplas políticas de evolução de schema

Não suporta evolução de schema

Preserva o changelog original

Interrompe a estrutura original do changelog

Lê e grava em várias tabelas

Lê e grava em uma única tabela

Em comparação com as instruções CTAS ou CDAS, os jobs Flink CDC oferecem as seguintes vantagens:

  • Sincroniza imediatamente a evolução de schema upstream, sem esperar que novas gravações de dados acionem o processo.

  • Preserva o changelog original, e as mensagens UPDATE não são divididas.

  • Sincroniza mais tipos de evolução de schema, como TRUNCATE TABLE e DROP TABLE.

  • Suporta mapeamentos de tabelas para definir nomes de tabelas de sink de forma flexível.

  • Permite comportamentos de evolução de schema flexíveis e configuráveis pelo usuário.

  • Possibilita filtragem de dados por meio de cláusulas WHERE.

  • Oferece suporte a pruning de colunas.

Flink CDC vs. Flink DataStream

A tabela a seguir apresenta as vantagens dos jobs de ingestão de dados Flink CDC em relação aos jobs DataStream.

Flink CDC

Flink DataStream

Projetado para usuários de todos os níveis de habilidade, não apenas especialistas

Exige conhecimento especializado em Java e sistemas distribuídos

Oculta detalhes de baixo nível e simplifica o desenvolvimento

Requer familiaridade com o framework Flink

O formato YAML é fácil de entender e aprender

Necessita de ferramentas como Maven para gerenciar dependências

Jobs existentes são fáceis de reutilizar

Código existente é difícil de reutilizar

Limitações

  • Use o Ververica Runtime (VVR) 11,1 ou superior para desenvolver jobs de ingestão de dados Flink CDC. Caso precise usar uma versão VVR 8.x, utilize a VVR 8.0.11.

  • Cada job suporta apenas uma source e um sink. Para ler de múltiplas sources ou gravar em múltiplos sinks, crie vários jobs Flink CDC.

  • Não implante jobs Flink CDC em um cluster de sessão.

  • Jobs Flink CDC não suportam ajuste automático.

Conectores de ingestão de dados Flink CDC

Consulte Conectores suportados para obter uma lista de conectores suportados como sources e sinks para ingestão de dados Flink CDC.

Criar um job de ingestão de dados Flink CDC

A partir de um modelo

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

  2. Na coluna Actions do workspace desejado, clique em Console.

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

  4. Clique em image e, em seguida, clique em New from Template.

  5. Selecione um modelo de sincronização de dados.

    Atualmente, existem modelos disponíveis para MySQL para StarRocks, MySQL para Paimon e MySQL para Hologres.

    image

  6. Insira as informações do job, incluindo nome do job, local de armazenamento e versão do mecanismo e, depois, clique em OK.

  7. Configure as informações de source e sink para o job Flink CDC.

    Consulte a documentação do conector correspondente para obter detalhes sobre a configuração de parâmetros.

A partir de CTAS/CDAS

Importante
  • Se um job contiver múltiplas instruções CTAS ou CDAS, o Flink detectará e converterá apenas a primeira.

  • Devido às diferenças no suporte a funções integradas entre Flink SQL e Flink CDC, as regras transform geradas podem não funcionar diretamente. Revise-as e ajuste-as conforme necessário.

  • Se a source for MySQL e o job CTAS/CDAS original ainda estiver em execução, ajuste o server-id da source do job de ingestão de dados Flink CDC para evitar conflitos.

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

  2. Na coluna Actions do workspace desejado, clique em Console.

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

  4. Clique em image e, em seguida, clique em New from CTAS/CDAS Job. Selecione o job CTAS ou CDAS desejado e clique em OK.

    A página de seleção exibe apenas jobs CTAS e CDAS válidos. Jobs ETL comuns ou rascunhos com erros de sintaxe não aparecem.

  5. Insira as informações do job, incluindo nome do job, local de armazenamento e versão do mecanismo e, depois, clique em OK.

A partir de open source

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

  2. Na coluna Actions do workspace desejado, clique em Console.

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

  4. Clique em image, selecione New Data Ingestion Draft, insira o File Name e a Engine Version e, em seguida, clique em Create.

  5. Copie o código do job Flink CDC open source.

  6. (Opcional) Clique em Validate.

    Valide a sintaxe, a conectividade de rede e as permissões de acesso.

Do zero

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

  2. Na coluna Actions do workspace desejado, clique em Console.

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

  4. Clique em image, selecione New Data Ingestion Draft, insira o File Name e a Engine Version e, em seguida, clique em Create.

  5. Configure o job Flink CDC.

    # Required
    source:
      # Source connector type
      type: <Replace with your source connector type>
      # Source configurations. For details, see the documentation for the corresponding connector.
      ...
    
    # Required
    sink:
      # Sink connector type
      type: <Replace with your sink connector type>
      # Sink configurations. For details, see the documentation for the corresponding connector.
      ...
    
    # Optional
    transform:
      # Transformation rule for the flink_test.customers table
      - source-table: flink_test.customers
        # Projection settings. Specifies the columns to synchronize and performs data transformation.
        projection: id, username, UPPER(username) as username1, age, (age + 1) as age1, test_col1, __schema_name__ || '.' || __table_name__ identifier_name
        # Filter condition. Synchronizes only data where id is greater than 10.
        filter: id > 10
        # Description of the transformation rule
        description: append calculated columns based on source table
    
    # Optional
    route:
      # Routing rule. Specifies the mapping between source tables and sink tables.
      - source-table: flink_test.customers
        sink-table: db.customers_o
        # Description of the routing rule
        description: sync customers table
      - source-table: flink_test.customers_suffix
        sink-table: db.customers_s
        # Description of the routing rule
        description: sync customers_suffix table
    
    # Optional
    pipeline:
      # Job name
      name: MySQL to Hologres Pipeline
    Nota

    Em um job Flink CDC, a chave e o valor devem ser separados por um espaço no formato Key: Value.

    A tabela a seguir descreve os blocos de código.

    Obrigatório

    Módulo

    Descrição

    Obrigatório

    source

    Início do pipeline de dados. O Flink CDC captura dados de alteração da source.

    Nota
    • Atualmente, o MySQL é a única source suportada. Para obter detalhes sobre itens de configuração específicos, consulte MySQL.

    • Use variáveis para gerenciar informações confidenciais. Para mais informações, consulte Gerenciamento de Variáveis.

    sink

    Fim do pipeline de dados. O Flink CDC transmite as alterações de dados capturadas para o sink.

    Nota

    Opcional

    pipeline

    (pipeline de dados)

    Define configurações básicas para todo o job de pipeline de dados, como o nome do pipeline.

    transform (transformação de dados)

    Especifica regras de transformação de dados. As transformações operam nos dados à medida que fluem pelo pipeline do Flink. Os recursos suportados incluem processamento ETL, filtragem com cláusulas WHERE, pruning de colunas e colunas calculadas.

    Use o módulo transform para transformar dados de alteração brutos do Flink CDC e adequá-los a um sistema downstream específico.

    route

    Se este módulo não for configurado, o job executará uma sincronização completa do banco de dados ou de tabelas especificadas.

    Em alguns casos, pode ser necessário enviar dados de alteração capturados para destinos diferentes com base em regras específicas. O mecanismo de roteamento permite especificar flexivelmente o relacionamento de mapeamento entre a source e o sink para enviar dados a sinks distintos.

    Para obter detalhes sobre a sintaxe e a configuração de cada módulo, consulte Referência para desenvolvimento de jobs de ingestão de dados Flink CDC.

    O código a seguir fornece um exemplo de como sincronizar todas as tabelas do banco de dados app_db no MySQL para um banco de dados no Hologres.

    source:
      type: mysql
      hostname: <hostname>
      port: 3306
      username: ${secret_values.mysqlusername}
      password: ${secret_values.mysqlpassword}
      tables: app_db.\.*
      server-id: 5400-5404
      # (Optional) Synchronize data from tables created during the incremental phase.
      scan.binlog.newly-added-table.enabled: true
      # (Optional) Synchronize table and column comments.
      include-comments.enabled: true
      # (Optional) Prioritize dispatching unbounded chunks to avoid potential TaskManager OutOfMemory issues.
      scan.incremental.snapshot.unbounded-chunk-first.enabled: true
      # (Optional) Enable parse filtering to speed up reads.
      scan.only.deserialize.captured.tables.changelog.enabled: true
    
    sink:
      type: hologres
      name: Hologres Sink
      endpoint: <endpoint>
      dbname: <database-name>
      username: ${secret_values.holousername}
      password: ${secret_values.holopassword}
    
    pipeline:
      name: Sync MySQL Database to Hologres
  6. (Opcional) Clique em Validate.

    Valide a sintaxe, a conectividade de rede e as permissões de acesso.

Documentos relacionados