Todos os produtos
Search
Central de documentação

AnalyticDB:Subscribe to binary logs by using Flink

Última atualização: Sep 10, 2026

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.

    Nota

    Para 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:

Etapa 1: Ativar o registro de binary log

  1. Ative o registro de binary log. Este exemplo usa uma tabela chamada source_table.

    Nota

    O 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;
  2. (Opcional) Modifique o período de retenção do binary log.

    Modifique o parâmetro binlog_ttl para 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_ttl aceita os seguintes formatos:

    • Milissegundos: valor numérico. Exemplo: 60 representa 60 milissegundos.

    • Segundos: número + s. Exemplo: 30s representa 30 segundos.

    • Horas: número seguido de h. Exemplo: 2h representa 2 horas.

    • Dias: número seguido de d. Exemplo: 1d representa 1 dia.

    Nota
    • O 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

  1. Baixe o arquivo JAR do conector: flink-sql-connector-adb-mysql-cdc-2,4-20260909.jar.

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

  3. Na aba Flink, localize o workspace desejado e clique em Console na coluna Actions.

  4. No painel de navegação à esquerda, clique em Connectors.

  5. Na página Connectors, clique em Create Custom Connector.

  6. Carregue o conector baixado e clique em Next.

  7. Clique em Finish. O conector personalizado criado aparece na lista de conectores.

Etapa 3: Assinar o binary log

  1. Faça login no console do Realtime Compute for Apache Flink e crie um job SQL.

  2. 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).

    Nota
    • A 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.file e scan.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.

    Importante

    Se 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.mode como 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.

  3. 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`)
    )
  4. 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.

  5. 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;
  6. Clique em Save.

  7. 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.

  8. (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.

  9. 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.

  10. (Opcional) Visualize as informações do binary log.

    Nota

    As 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