A instrução CREATE DATABASE AS (CDAS) sincroniza schemas e dados de tabelas em todo um banco de dados em tempo real, incluindo alterações de schema. Use-a para replicar várias — ou todas — as tabelas de um banco de dados de origem para um destino sem criar manualmente as tabelas de sink antecipadamente.
Para novos jobs, use YAML para ingestão de dados . O YAML oferece suporte a todos os recursos principais do CDAS — sincronização de banco de dados e tabelas, evolução de schema, colunas computadas personalizadas, sincronização de log binário bruto, filtragem por cláusula WHERE e pruning de colunas — e permite converter rascunhos SQL existentes contendo instruções CTAS ou CDAS . Para mais detalhes, consulte Ingestão de dados com Flink CDC .
Como funciona
O CDAS é um recurso sintático do CREATE TABLE AS (CTAS). Ao executar uma instrução CDAS, o Realtime Compute for Apache Flink:
-
Verifica se o banco de dados de destino e as tabelas de sink existem.
Se o banco de dados de destino não existir, o Flink o criará por meio do catálogo de destino.
-
Se já existir, o Flink verificará as tabelas de sink:
Na ausência de tabelas de sink, o Flink as criará com nomes e schemas idênticos aos das tabelas de origem.
Se as tabelas de sink já estiverem presentes, o Flink ignorará a criação.
Inicia o job de sincronização de dados. O sistema replica continuamente dados e alterações de schema do banco de dados de origem para o destino.
Por que usar CDAS em vez de múltiplas instruções CTAS:
Sintaxe simplificada: O Flink expande automaticamente uma única instrução CDAS em uma instrução CTAS para cada tabela.
Uso otimizado de recursos: O Flink utiliza um único vértice de source para ler todas as tabelas correspondentes. Em fontes MySQL Change Data Capture (CDC), isso reduz conexões com o banco de dados, evita leituras redundantes de log binário e diminui a carga geral de leitura no banco de dados MySQL.
Recursos principais
Sincronização de dados
|
Recurso |
Descrição |
|
Sincroniza dados completa e incrementalmente de várias tabelas (ou de todas) em um banco de dados para cada tabela de sink relacionada. |
|
|
Identifica nomes de tabelas de origem entre shards de banco de dados usando expressões regulares, consolida essas tabelas e as sincroniza com os sinks correspondentes. |
|
|
Sincroniza tabelas recém-adicionadas ao reiniciar o job a partir de um savepoint. |
|
|
Permite usar a instrução STATEMENT SET para submeter múltiplas instruções CDAS e CTAS como um único job. Também possibilita mesclar e reutilizar dados dos operadores de tabela de origem para reduzir a carga de leitura na fonte de dados. |
Evolução de schema
Durante a sincronização de um banco de dados via CDAS, alterações de schema — como adição de colunas — são propagadas automaticamente para o sink. O comportamento e a política são os mesmos do CTAS. Consulte Evolução de schema.
Pré-requisitos
Antes de começar, certifique-se de ter:
Um catálogo de destino registrado no workspace. Consulte Catálogos.
Limitações
Limitações de sintaxe
Não há suporte para depuração de rascunhos SQL que contenham uma instrução CDAS.
-
MiniBatch não é suportado.
ImportanteRemova todas as configurações do MiniBatch antes de criar um rascunho SQL com uma instrução CTAS ou CDAS: 1. Acesse O&M > Configurations. 2. Selecione a aba Deployment Defaults. 3. Na seção Other Configuration, verifique se as configurações do MiniBatch foram removidas. Se encontrar o erro "Currently does not support merge StreamExecMiniBatchAssigner type ExecNode in CTAS/CDAS syntax" ao criar ou iniciar um deployment, consulte Como corrigir este erro?
Conectores suportados
|
Conector |
Source |
Sink |
Observações |
|
Sim |
Não |
Views não podem ser sincronizadas. |
|
|
Sim |
Não |
— |
|
|
Sim |
Não |
Não há suporte para consolidação de tabelas e bancos de dados com sharding. Metadados do MongoDB não podem ser sincronizados. Para configuração, consulte Gerenciar catálogos do MongoDB. |
|
|
Não |
Sim |
— |
|
|
Não |
Sim |
O Flink cria uma conexão para cada tabela com base na opção |
|
|
Não |
Sim |
Limitado ao StarRocks no Alibaba Cloud EMR. |
|
|
Não |
Sim |
A sincronização de dados para tabelas de sink Paimon do metastore DLF 2.5 é suportada no VVR 11.1 ou posterior. |
Notas de uso
Sincronização de novas tabelas
VVR 8.0.6 ou posterior: Reinicie o job a partir de um savepoint após adicionar uma nova tabela. Consulte Sincronizar novas tabelas.
-
VVR 8.0.5 ou anterior: Reinícios de job não capturam novas tabelas. Utilize uma das seguintes abordagens:
Abordagem
Passos
Criar um novo job para a nova tabela
Mantenha o job existente em execução. Crie um job separado direcionado apenas à nova tabela. Exemplo:
CREATE TABLE IF NOT EXISTS new_table AS TABLE mysql.tpcds.new_table /*+ OPTIONS('server-id'='8008-8010') */;Limpar e reiniciar do zero
1. Cancele o job existente. 2. Limpe os dados sincronizados no sink. 3. Reinicie o job sem estados.
Acesso entre contas e RAM
Conceda as permissões de leitura e gravação necessárias à sua conta ao acessar recursos externos entre contas ou como usuário RAM ou função RAM.
Sintaxe
CREATE DATABASE IF NOT EXISTS <target_database>
[COMMENT database_comment]
[WITH (key1=val1, key2=val2, ...)]
AS DATABASE <source_database>
INCLUDING { ALL TABLES | TABLE 'table_name' }
[EXCLUDING TABLE 'table_name']
[/*+ OPTIONS(key1=val1, key2=val2, ... ) */]
<target_database>:
[catalog_name.]db_name
<source_database>:
[catalog_name.]db_name
IF NOT EXISTS é obrigatório. Isso instrui o Flink a verificar se a tabela de sink existe no destino. Se estiver ausente, o Flink a cria; se já existir, o Flink ignora a criação.
A tabela de sink compartilha o schema da tabela de origem — chave primária e nomes e tipos de campos físicos — mas exclui colunas computadas, campos de metadados e configurações de watermark.
O Flink mapeia tipos de dados da origem para o sink. Para detalhes, consulte a documentação do conector relevante.
Parâmetros
|
Parâmetro |
Descrição |
|
|
|
Nome do banco de dados de destino, opcionalmente prefixado com o nome do catálogo: |
|
|
|
Descrição para o banco de dados de destino. O padrão é a descrição do banco de dados de origem. |
|
|
|
Opções para o banco de dados de destino. Chaves e valores devem ser strings, por exemplo, |
|
|
|
Nome do banco de dados de origem, opcionalmente prefixado com o nome do catálogo: |
|
|
|
Sincroniza todas as tabelas no banco de dados de origem. |
|
|
|
Especifica as tabelas a serem sincronizadas. Separe múltiplas tabelas com |
`. Suporta expressões regulares — por exemplo, |
|
|
Especifica as tabelas a serem excluídas da sincronização. Separe múltiplas tabelas com |
`. Suporta expressões regulares — por exemplo, |
|
|
Opções de conector para a tabela de origem. Chaves e valores devem ser strings, por exemplo, |
Exemplos
Sincronizar um banco de dados
Sincronize todas as tabelas do banco de dados MySQL tpcds para o Hologres.
Pré-requisitos:
Um catálogo Hologres chamado
holodeve estar criado no workspace.Um catálogo MySQL chamado
mysqldeve estar criado no workspace.
USE CATALOG holo;
CREATE DATABASE IF NOT EXISTS holo_tpcds -- Create holo_tpcds in Hologres.
WITH ('sink.parallelism' = '4') -- Set sink parallelism (default: 4 for Hologres).
AS DATABASE mysql.tpcds INCLUDING ALL TABLES -- Sync all tables from mysql.tpcds.
/*+ OPTIONS('server-id'='8001-8004') */; -- Configure the MySQL CDC source options.
As opções definidas na cláusula WITH aplicam-se apenas ao job atual e controlam o comportamento de escrita. Elas não são persistidas no catálogo Hologres. Para opções suportadas, consulte Hologres .
Consolidar e sincronizar shards de banco de dados
Mescle todas as tabelas com nomes idênticos em vários shards de banco de dados MySQL em tabelas únicas no Hologres.
Cenário: Uma instância MySQL possui shards chamados order_db01 até order_db99. Cada shard contém tabelas como order e order_detail. Todos os dados dos shards — incluindo alterações de schema — precisam ser enviados ao Hologres.
Solução: Use uma expressão regular no nome do banco de dados para corresponder a todos os shards. Os nomes do banco de dados e da tabela são adicionados a cada tabela de sink como dois campos adicionais. A chave primária do Hologres inclui o nome do banco de dados, o nome da tabela e as colunas de chave primária da tabela de origem para garantir unicidade. Não é necessário criar tabelas de destino antecipadamente.
USE CATALOG holo;
CREATE DATABASE IF NOT EXISTS holo_order -- Create holo_order in Hologres.
WITH('sink.parallelism'='4') -- Set sink parallelism (optional).
AS DATABASE mysql.`order_db[0-9]+` INCLUDING ALL TABLES -- Match all order_db shards.
/*+ OPTIONS('server-id'='8001-8004') */; -- Configure MySQL CDC source options (optional).
Tabelas com nomes idênticos entre shards são mescladas em uma única tabela Hologres.
Sincronizar novas tabelas
Após o início de um job CDAS em execução, ative a detecção de novas tabelas e reinicie a partir de um savepoint para capturar tabelas recém-adicionadas.
A detecção de novas tabelas requer VVR 8.0.6 ou posterior. O modo de inicialização da tabela de origem deve ser initial .
Na página Deployments, localize o deployment desejado e clique em Cancel na coluna Actions.
Na caixa de diálogo, expanda More Strategies, selecione Stop With Savepoint e clique em OK.
-
No rascunho SQL, adicione a seguinte instrução:
SET 'table.cdas.scan.newly-added-table.enabled' = 'true'; Clique em Deploy.
-
Recupere o job a partir do savepoint:
Na página Deployments, clique no nome do deployment.
Na página de detalhes do deployment, clique na aba State e depois na subaba History.
Na lista Savepoints, localize o savepoint criado quando você parou o job.
Escolha More > Start job from this savepoint na coluna Actions. Consulte Iniciar um job.
Executar múltiplas instruções CDAS e CTAS como um único job
Sincronize vários bancos de dados em um único job para maximizar a reutilização da source.
Cenário: Sincronizar tpcds, tpch e os shards user_db01–user_db99 para o Hologres em um único job.
Solução: Envolva todas as instruções em um bloco STATEMENT SET. O Flink reutiliza um único vértice de source em todas as instruções, reduzindo o número de server IDs, conexões de banco de dados e a carga geral de leitura — especialmente para fontes MySQL CDC.
As opções do conector devem ser idênticas em todas as tabelas de origem para permitir a reutilização da source.
Para configuração de server ID, consulte Definir o server ID para evitar conflitos de consumo de binlog.
USE CATALOG holo;
BEGIN STATEMENT SET;
-- Sync user tables from all user_db shards.
CREATE TABLE IF NOT EXISTS user
AS TABLE mysql.`user_db[0-9]+`.`user[0-9]+`
/*+ OPTIONS('server-id'='8001-8004') */;
-- Sync tpcds database.
CREATE DATABASE IF NOT EXISTS holo_tpcds
AS DATABASE mysql.tpcds INCLUDING ALL TABLES
/*+ OPTIONS('server-id'='8001-8004') */;
-- Sync tpch database.
CREATE DATABASE IF NOT EXISTS holo_tpch
AS DATABASE mysql.tpch INCLUDING ALL TABLES
/*+ OPTIONS('server-id'='8001-8004') */;
END;
Sincronizar múltiplos bancos de dados MySQL para o Kafka
Sincronize tabelas de vários bancos de dados MySQL para o Kafka evitando conflitos de nomes de tópicos.
Cenário: Tabelas com o mesmo nome existem tanto em tpcds quanto em tpch. Enviar ambas para o mesmo tópico Kafka sobrescreveria os dados.
Solução: Utilize a opção cdas.topic.pattern para gerar nomes de tópicos exclusivos por banco de dados. O placeholder {table-name} é substituído em tempo de execução pelo nome real da tabela. Por exemplo, 'cdas.topic.pattern'='tpcds-{table-name}' direciona a table1 do tpcds para o tópico tpcds-table1.
USE CATALOG kafkaCatalog;
BEGIN STATEMENT SET;
-- Sync tpcds database. Topic names follow the pattern "tpcds-<table_name>".
CREATE DATABASE IF NOT EXISTS kafka
WITH ('cdas.topic.pattern' = 'tpcds-{table-name}')
AS DATABASE mysql.tpcds INCLUDING ALL TABLES
/*+ OPTIONS('server-id'='8001-8004') */;
-- Sync tpch database. Topic names follow the pattern "tpch-<table_name>".
CREATE DATABASE IF NOT EXISTS kafka
WITH ('cdas.topic.pattern' = 'tpch-{table-name}')
AS DATABASE mysql.tpch INCLUDING ALL TABLES
/*+ OPTIONS('server-id'='8001-8004') */;
END;
A introdução do Kafka como camada intermediária entre o MySQL e o Flink também reduz a carga no MySQL. Consulte Sincronizar banco de dados MySQL para Kafka com Flink CDC.
FAQ
Erros de execução
Performance do job
Sincronização de dados
Próximos passos
-
Catálogos comumente usados com CDAS:
-
Recursos relacionados:
-
Ingestão de dados via YAML: