O Realtime Compute for Apache Flink assina os binary logs de um cluster do AnalyticDB for MySQL para capturar e processar alterações no banco de dados em tempo real. Isso permite sincronização eficiente de dados e computação em fluxo.
Pré-requisitos
O cluster do AnalyticDB for MySQL deve ser uma das seguintes edições: Enterprise Edition, Basic Edition, Data Lakehouse Edition ou Data Warehouse Edition (in elastic mode).
-
O cluster do AnalyticDB for MySQL deve atender aos seguintes requisitos de versão do kernel:
Mecanismo de tabela xuanwu_v1: 3.2.1.0 ou posterior.
Mecanismo de tabela xuanwu_v2: 3.2.5.0 ou posterior.
Para assinar binary logs de visualizações materializadas incrementais, a versão do kernel deve ser 3.2.6.9, 3.2.7.1 ou posterior.
NotaPara visualize e atualize a versão secundária, acesse a seção Configuration Information na página Cluster Information no console do AnalyticDB for MySQL.
O workspace do Flink deve usar o Ververica Runtime (VVR) 8.0.4 ou versão posterior.
O cluster do AnalyticDB for MySQL e o fully managed Flink workspace devem estar na mesma VPC.
Adicione o CIDR block do workspace do Flink à whitelist do AnalyticDB for MySQL.
Limitações
O Flink processa apenas basic data types e o JSON complex data type do binary log do AnalyticDB for MySQL.
-
As operações a seguir não são capturadas ao usar o Flink para assinar binary logs do AnalyticDB for MySQL. As alterações de dados correspondentes não são sincronizadas com os consumidores downstream:
Operações DDL: Alterações de esquema, como ALTER TABLE, não são capturadas. Se o esquema da tabela de origem mudar, atualize manualmente a definição da tabela no job do Flink e reimplante o job.
INSERT OVERWRITE: Dados gravados por meio de INSERT OVERWRITE não geram registros de binary log. Para obter os dados completos após uma operação INSERT OVERWRITE, reassine os binary logs.
Evicção automática de partições: Quando partições expiradas em uma tabela particionada são removidas automaticamente, os dados evictos não geram registros de binary log.
-
Atualização completa de visualizações materializadas: Operações de atualização completa em visualizações materializadas não geram registros de binary log.
NotaOs binary logs de visualizações materializadas contêm apenas alterações de dados provenientes de atualizações incrementais. Para obter o estado completo dos dados após uma atualização completa, reassine os binary logs.
Etapa 1: Ativar o registro de binary log
-
Ative o registro de binary log. Este exemplo usa uma tabela chamada source_table.
NotaO AnalyticDB for MySQL suporta a ativação de binary log apenas no nível da tabela.
Ao criar uma tabela
CREATE TABLE source_table ( `id` INT, `num` BIGINT, PRIMARY KEY (`id`) )DISTRIBUTED BY HASH (id) BINLOG=true;Para uma tabela existente
ALTER TABLE source_table BINLOG=true; -
(Opcional) Modifique o período de retenção do binary log.
Modifique o parâmetro
binlog_ttlpara ajustar o período de retenção do binary log. O valor padrão desse parâmetro é 6 horas. Por exemplo, defina o período de retenção do binary log da tabela source_table como 1 dia.ALTER TABLE source_table binlog_ttl='1d';O parâmetro
binlog_ttlaceita os seguintes formatos:Milissegundos: valor numérico. Exemplo:
60representa 60 milissegundos.Segundos: número + s. Exemplo:
30srepresenta 30 segundos.Horas: número seguido de h. Exemplo:
2hrepresenta 2 horas.Dias: número seguido de
d. Exemplo:1drepresenta 1 dia.
NotaO período máximo de retenção do binary log é de 365 dias para clusters com as seguintes versões de kernel ou posteriores em suas respectivas séries: 3.2.1.9, 3.2.2.14, 3.2.3.8, 3.2.4.4 e 3.2.5.1. Para clusters com versões anteriores do kernel, o período máximo de retenção é de 21 dias.
Recomendamos definir o período de retenção do binary log com um valor não inferior ao valor padrão do parâmetro
binlog_ttl. Se o período de retenção for muito curto, os arquivos podem ser removidos, o que afeta a sincronização de dados.Para visualizar o período atual de retenção do binary log, execute
SHOW CREATE TABLE source_table;.
Etapa 2: Carregar o conector do AnalyticDB for MySQL
Baixe o arquivo JAR do conector: flink-sql-connector-adb-mysql-cdc-2,4-20260909.jar.
Faça login no console do Realtime Compute for Apache Flink.
Na aba Flink, localize o workspace desejado e clique em Console na coluna Actions.
No painel de navegação à esquerda, clique em Connectors.
Na página Connectors, clique em Create Custom Connector.
Carregue o conector baixado e clique em Next.
Clique em Finish. O conector personalizado criado aparece na lista de conectores.
Etapa 3: Assinar o binary log
Faça login no console do Realtime Compute for Apache Flink e crie um job SQL.
-
Crie uma tabela de origem que se conecte ao AnalyticDB for MySQL e leia dados de binary log de uma tabela específica (source_table).
NotaA chave primária definida na DDL do Flink deve corresponder à chave primária da tabela física no cluster do AnalyticDB for MySQL. Isso inclui tanto as colunas de chave quanto o nome da chave primária. Caso contrário, os dados podem ficar incorretos.
Os tipos de dados do Flink devem ser compatíveis com os do AnalyticDB for MySQL. Para detalhes sobre os mapeamentos de tipos de dados, consulte Type mapping.
CREATE TEMPORARY TABLE adb_source ( `id` INT, `num` BIGINT, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'adb-mysql-cdc', 'hostname' = 'amv-2zepb9n1l58ct01z50000****.ads.aliyuncs.com', 'username' = 'testUser', 'password' = 'Test12****', 'database-name' = 'binlog', 'table-name' = 'source_table' );A tabela a seguir descreve os parâmetros WITH.
Parâmetro
Obrigatório
Padrão
Tipo
Descrição
connector
Sim
Nenhum
STRING
Conector a ser usado.
Este é um conector personalizado. Defina este parâmetro como
adb-mysql-cdc.hostname
Sim
Nenhum
STRING
Endpoint do cluster do AnalyticDB for MySQL.
username
Sim
Nenhum
STRING
Conta do banco de dados do cluster do AnalyticDB for MySQL.
password
Sim
Nenhum
STRING
Senha da conta do banco de dados do AnalyticDB for MySQL.
database-name
Sim
Nenhum
STRING
Nome do banco de dados do AnalyticDB for MySQL.
Como o AnalyticDB for MySQL implementa registro de binary log no nível da tabela, é possível especificar apenas um banco de dados.
table-name
Sim
Nenhum
STRING
Nome da tabela no banco de dados do AnalyticDB for MySQL.
Como o AnalyticDB for MySQL implementa registro de binary log no nível da tabela, é possível especificar apenas uma tabela.
port
Não
3306
INTEGER
Número da porta.
scan.incremental.snapshot.enabled
Não
true
BOOLEAN
Especifica se o mecanismo de leitura de snapshot incremental deve ser ativado.
Esse recurso está ativado por padrão. Um snapshot incremental é um novo mecanismo para leitura de snapshots de tabelas. Em comparação com o mecanismo anterior, os snapshots incrementais oferecem os seguintes benefícios:
A origem suporta leitura simultânea durante a fase de snapshot.
A origem suporta checkpoints no nível de chunk durante a leitura do snapshot.
A origem não precisa obter permissões de bloqueio do banco de dados antes de ler um snapshot.
scan.incremental.snapshot.chunk.size
Não
8096
INTEGER
Número de linhas por chunk para um snapshot de tabela.
scan.snapshot.fetch.size
Não
1024
INTEGER
Número máximo de linhas a serem buscadas por vez durante a leitura de um snapshot de tabela.
scan.startup.mode
Não
initial
STRING
Modo de inicialização para consumo de dados.
Valores válidos:
initial (padrão): executa um snapshot inicial da tabela e, em seguida, lê o binary log mais recente.
earliest-offset: ignora a fase de snapshot e começa a leitura a partir do binary log mais antigo disponível.
specific-offset: ignora a fase de snapshot e inicia a partir de uma posição específica do binary log. Especifique o nome do arquivo e o deslocamento do binary log definindo os parâmetros
scan.startup.specific-offset.fileescan.startup.specific-offset.pos.latest-offset: ignora a fase de snapshot e lê apenas as alterações ocorridas após a inicialização do conector.
timestamp: ignora a fase de snapshot e começa a leitura a partir de um carimbo de data/hora especificado. Defina o carimbo de data/hora em milissegundos (ms) usando o parâmetro
scan.startup.timestamp-millis.
ImportanteSe você usar o modo de inicialização earliest-offset, specific-offset ou timestamp, garanta que o esquema da tabela permaneça inalterado entre a posição inicial especificada e a inicialização do job. Caso contrário, o job pode falhar devido a alterações no esquema.
scan.startup.specific-offset.file
Não
Nenhum
STRING
No modo de inicialização specific-offset, indica o nome do arquivo de binary log na posição inicial.
Para obter o nome do arquivo de binary log mais recente, execute a instrução
SHOW MASTER STATUS for table_name;.scan.startup.specific-offset.pos
Não
Nenhum
LONG
No modo de inicialização specific-offset, indica a posição no arquivo de binary log na posição inicial.
Para obter a posição mais recente do binary log, execute a instrução
SHOW MASTER STATUS for table_name;.scan.startup.specific-offset.skip-events
Não
Nenhum
LONG
Número de eventos a serem ignorados após a posição inicial especificada.
scan.startup.specific-offset.skip-rows
Não
Nenhum
LONG
Número de linhas de dados a serem ignoradas após a posição inicial especificada.
scan.startup.timestamp-millis
Não
Nenhum
LONG
Carimbo de data/hora em milissegundos da posição inicial ao usar o modo de inicialização timestamp.
Ao usar este parâmetro, defina
scan.startup.modecomo timestamp. O carimbo de data/hora está em milissegundos (ms).server-time-zone
Não
Padrão do sistema
STRING
Fuso horário da sessão no servidor de banco de dados, como "Asia/Shanghai".
Controla como o AnalyticDB for MySQL converte valores para o tipo de dados
STRING. Se este parâmetro não for definido,ZONELD.SYSTEMDEFAULT()será usado para determinar o fuso horário do servidor.debezium.min.row.count.to.stream.result
Não
1000
INTEGER
Se o número de linhas em uma tabela for maior que esse valor, o conector transmitirá os resultados em fluxo.
Se você definir este parâmetro como
0, todas as verificações de tamanho da tabela serão ignoradas e todos os resultados serão transmitidos em fluxo durante a fase de snapshot.connect.timeout
Não
30s
DURATION
Tempo máximo de espera por uma conexão de banco de dados antes que a tentativa expire.
A unidade padrão é segundos (s).
connect.max-retries
Não
3
INTEGER
Número máximo de tentativas após uma falha na conexão com o banco de dados.
-
Crie uma tabela física no banco de dados de destino para armazenar os dados processados. Este tópico usa o AnalyticDB for MySQL como destino. Para conectores suportados pelo Flink, consulte Supported connectors.
CREATE TABLE target_table ( `id` INT, `num` BIGINT, PRIMARY KEY (`id`) ) -
Crie uma tabela sink que se conecte à tabela criada na etapa anterior. A tabela sink grava os dados processados na tabela especificada no AnalyticDB for MySQL.
CREATE TEMPORARY TABLE adb_sink ( `id` INT, `num` BIGINT, PRIMARY KEY (`id`) NOT ENFORCED ) WITH ( 'connector' = 'adb3.0', 'url' = 'jdbc:mysql://amv-2zepb9n1l58ct01z50000****.ads.aliyuncs.com:3306/flinktest', 'userName' = 'testUser', 'password' = 'Test12****', 'tableName' = 'target_table' );Para mais informações sobre os parâmetros WITH e mapeamentos de tipos para a tabela sink, consulte AnalyticDB for MySQL V3.0 connector.
-
Use uma instrução INSERT INTO para enviar dados da tabela de origem para a tabela sink.
INSERT INTO adb_sink SELECT * FROM adb_source; Clique em Save.
-
Clique em Validate.
O recurso de validação verifica a semântica SQL, a conectividade de rede e os metadados das tabelas usadas pelo job. Você também pode clicar em SQL Advice na área de resultados para visualizar alertas de risco SQL e sugestões de otimização.
-
(Opcional) Clique em Debug.
Use o recurso de job debugging para simular execuções de jobs, verificar resultados de saída, validar a lógica de negócios das instruções SELECT ou INSERT, aumentar a eficiência de desenvolvimento e reduzir riscos de qualidade de dados.
-
Clique em Deploy.
Após desenvolver e validar o job, deploy it para o ambiente de produção. Em seguida, acesse a página O&M e start the job.
-
(Opcional) Visualize as informações do binary log.
NotaAs instruções a seguir retornam 0 se você tiver ativado o registro de binary log, mas ainda não tiver assinado usando o DTS. As informações do binary log aparecem somente após o estabelecimento de uma assinatura bem-sucedida.
-
Para obter o nome do arquivo e a posição da entrada mais recente do binary log, execute a seguinte instrução SQL:
SHOW MASTER STATUS FOR source_table; -
Para visualizar todos os arquivos de binary log não removidos e seus tamanhos, execute a seguinte instrução SQL:
SHOW BINARY LOGS FOR source_table;
-
Mapeamento de tipos
A tabela a seguir mapeia os tipos de dados do AnalyticDB for MySQL para seus equivalentes no Flink.
|
**Tipo do AnalyticDB for MySQL** |
**Tipo do Flink** |
|
BOOLEAN |
BOOLEAN |
|
TINYINT |
TINYINT |
|
SMALLINT |
SMALLINT |
|
INT |
INT |
|
BIGINT |
BIGINT |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
DECIMAL(p,s) ou NUMERIC(p,s) |
DECIMAL(p,s) |
|
VARCHAR |
STRING |
|
BINARY |
BYTES |
|
DATE |
DATE |
|
TIME |
TIME |
|
DATETIME |
TIMESTAMP |
|
TIMESTAMP |
TIMESTAMP |
|
POINT |
STRING |
|
JSON |
STRING |