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
SchemaChangeEventpara transmitir informações de evolução de schema, como criação de tabela, adição de coluna ou truncamento de tabela. Eles utilizamDataChangeEventpara 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 |
|
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
UPDATEnão são divididas.Sincroniza mais tipos de evolução de schema, como
TRUNCATE TABLEeDROP 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
Faça login no console do Realtime Compute for Apache Flink.
Na coluna Actions do workspace desejado, clique em Console.
No painel de navegação à esquerda, escolha .
Clique em
e, em seguida, clique em New from Template.-
Selecione um modelo de sincronização de dados.
Atualmente, existem modelos disponíveis para MySQL para StarRocks, MySQL para Paimon e MySQL para Hologres.

Insira as informações do job, incluindo nome do job, local de armazenamento e versão do mecanismo e, depois, clique em OK.
-
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
Se um job contiver múltiplas instruções
CTASouCDAS, 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
transformgeradas podem não funcionar diretamente. Revise-as e ajuste-as conforme necessário.Se a source for MySQL e o job
CTAS/CDASoriginal ainda estiver em execução, ajuste oserver-idda source do job de ingestão de dados Flink CDC para evitar conflitos.
Faça login no console do Realtime Compute for Apache Flink.
Na coluna Actions do workspace desejado, clique em Console.
No painel de navegação à esquerda, escolha .
-
Clique em
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.
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
Faça login no console do Realtime Compute for Apache Flink.
Na coluna Actions do workspace desejado, clique em Console.
No painel de navegação à esquerda, escolha .
Clique em
, selecione New Data Ingestion Draft, insira o File Name e a Engine Version e, em seguida, clique em Create.Copie o código do job Flink CDC open source.
-
(Opcional) Clique em Validate.
Valide a sintaxe, a conectividade de rede e as permissões de acesso.
Do zero
Faça login no console do Realtime Compute for Apache Flink.
Na coluna Actions do workspace desejado, clique em Console.
No painel de navegação à esquerda, escolha .
Clique em
, selecione New Data Ingestion Draft, insira o File Name e a Engine Version e, em seguida, clique em Create.-
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 PipelineNotaEm 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.
NotaAtualmente, 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.
NotaPara ver os sistemas de sink suportados, consulte Conectores de ingestão de dados Flink CDC. Para obter detalhes sobre os itens de configuração do sink, consulte a documentação do conector correspondente.
Use variáveis para gerenciar informações confidenciais. Para mais informações, consulte Gerenciamento de Variáveis.
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
transformpara 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 -
(Opcional) Clique em Validate.
Valide a sintaxe, a conectividade de rede e as permissões de acesso.
Documentos relacionados
Após desenvolver um job Flink CDC, implante-o. Consulte Implantar um job para obter instruções de implantação.
Para criar rapidamente um job Flink CDC que sincroniza dados de um banco de dados MySQL para o StarRocks, consulte Início rápido para jobs de ingestão de dados Flink CDC.