Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector MySQL YAML

Última atualização: Sep 18, 2026

O conector MySQL pode ser usado como fonte de dados em jobs YAML de ingestão de dados.

Pré-requisitos

Antes de usar uma tabela de origem MySQL CDC, conclua as operações de pré-requisito descritas em Configure MySQL.

ApsaraDB RDS for MySQL

  • Execute uma network probe para garantir a conectividade de rede com o Realtime Compute for Apache Flink.

  • Versão do MySQL: 5,6, 5,7, 8.0.x ou 8,4.

  • O log binário deve estar ativado. Essa configuração vem habilitada por padrão.

  • O formato do log binário deve ser ROW. Este é o formato padrão.

  • O parâmetro binlog_row_image deve estar definido como FULL. Esta é a configuração padrão.

  • A compactação de transações de log binário deve estar desativada. Esse recurso foi introduzido no MySQL 8.0.20 e vem desabilitado por padrão.

  • Crie um usuário MySQL com as permissões SELECT, SHOW DATABASES, REPLICATION SLAVE e REPLICATION CLIENT.

  • Crie um banco de dados e uma tabela MySQL. Para mais informações, consulte Create a database and an account for an ApsaraDB RDS for MySQL instance. Use uma conta privilegiada para criar o banco de dados MySQL e evitar falhas na operação devido a permissões insuficientes.

  • Configure uma lista de permissões de endereços IP. Para mais informações, consulte Configure an IP address whitelist for an ApsaraDB RDS for MySQL instance.

PolarDB for MySQL

  • Execute uma network probe para garantir a conectividade de rede com o Realtime Compute for Apache Flink.

  • Versão do MySQL: 5,6, 5,7, 8.0.x ou 8,4.

  • O log binário deve estar ativado. Essa configuração vem desabilitada por padrão.

  • O formato do log binário deve ser ROW. Este é o formato padrão.

  • O parâmetro binlog_row_image deve estar definido como FULL. Esta é a configuração padrão.

  • A compactação de transações de log binário deve estar desativada. Esse recurso foi introduzido no MySQL 8.0.20 e vem desabilitado por padrão.

  • Crie um usuário MySQL com as permissões SELECT, SHOW DATABASES, REPLICATION SLAVE e REPLICATION CLIENT.

  • Crie um banco de dados e uma tabela MySQL. Para mais informações, consulte Create a database and an account for a PolarDB for MySQL cluster. Use uma conta privilegiada para criar o banco de dados MySQL e evitar falhas na operação devido a permissões insuficientes.

  • Configure uma lista de permissões de endereços IP. Para mais informações, consulte Configure an IP address whitelist for a PolarDB for MySQL cluster.

Self-managed MySQL

  • Execute uma network probe para garantir a conectividade de rede com o Realtime Compute for Apache Flink.

  • Versão do MySQL: 5,6, 5,7, 8.0.x ou 8,4.

  • O log binário deve estar ativado. Essa configuração vem desabilitada por padrão.

  • O formato do log binário deve ser ROW. O formato padrão é STATEMENT.

  • O parâmetro binlog_row_image deve estar definido como FULL. Esta é a configuração padrão.

  • A compactação de transações de log binário deve estar desativada. Esse recurso foi introduzido no MySQL 8.0.20 e vem desabilitado por padrão.

  • Crie um usuário MySQL e conceda as permissões SELECT, SHOW DATABASES, REPLICATION SLAVE e REPLICATION CLIENT.

  • Crie um banco de dados e uma tabela MySQL. Para mais informações, consulte Create a database and an account for a self-managed MySQL instance. Use uma conta privilegiada para criar o banco de dados MySQL e evitar falhas na operação devido a permissões insuficientes.

  • Configure uma lista de permissões de endereços IP. Para mais informações, consulte Configure an IP address whitelist for a self-managed MySQL instance.

Limites

Limites gerais

O conector MySQL CDC não oferece suporte ao recurso de compactação de transações de log binário. Portanto, ao usar o conector MySQL CDC para consumir dados incrementais, certifique-se de que a compactação de transações de log binário esteja desativada. Caso contrário, o conector poderá falhar ao recuperar dados incrementais.

Limites do ApsaraDB RDS for MySQL

  • No ApsaraDB RDS for MySQL, não leia dados de um banco de dados secundário ou réplica somente leitura. O período de retenção padrão de logs binários para bancos secundários e réplicas somente leitura é curto. Se os logs binários expirarem e forem limpos, o job poderá falhar ao consumir os dados do log binário e gerar um erro.

  • O ApsaraDB RDS for MySQL habilita a sincronização paralela primária/secundária por padrão, mas não garante uma ordem consistente de transações entre as instâncias primária e secundária. Isso pode causar perda de dados durante um failover primário/secundário e recuperação de checkpoint. Para evitar esse problema, habilite manualmente a opção slave_preserve_commit_order no ApsaraDB RDS for MySQL.

Limites do PolarDB for MySQL

As tabelas de origem MySQL CDC não suportam a leitura de dados de clusters com arquitetura Multi-master Cluster do PolarDB for MySQL V1.0.19 e versões anteriores. Para mais informações, consulte What is a Multi-master Cluster?. Os logs binários gerados por esses clusters podem conter IDs de tabela duplicados, causando erros de mapeamento de esquema na tabela de origem CDC e falhas ao analisar os dados do log binário.

Limites do MySQL open source

Por padrão, o MySQL mantém a ordem das transações durante a replicação de log binário primária/secundária. Se uma réplica MySQL tiver a replicação paralela habilitada (slave_parallel_workers > 1), mas não tiver a opção slave_preserve_commit_order=ON ativada, a ordem de commit das transações poderá ficar inconsistente em relação ao banco de dados primário. Quando o Flink CDC se recupera de um checkpoint, ele pode perder dados devido à sequência desordenada. Defina slave_preserve_commit_order = ON na réplica MySQL. Alternativamente, defina slave_parallel_workers = 1, mas isso sacrificará o desempenho da replicação.

Observações de uso

  • Set a server ID to avoid binary log consumption conflicts.

  • Durante a fase de leitura completa de dados, não é possível salvar um savepoint, adicionar ou remover uma tabela da tabela de origem e reiniciar o job a partir desse savepoint. Essa ação causa falha na leitura dos dados.

Ingestão de dados

Use o conector MySQL como fonte de dados em um job YAML de ingestão de dados.

Sintaxe

source:
   type: mysql
   name: MySQL Source
   hostname: localhost
   port: 3306
   username: <username>
   password: <password>
   tables: adb.\.*, bdb.user_table_[0-9]+, [app|web].order_\.*
   server-id: 5401-5404

sink:
  type: xxx

Itens de configuração

Parâmetro

Descrição

Obrigatório

Tipo de dado

Valor padrão

Observações

type

O tipo da fonte de dados.

Sim

STRING

Nenhum

O valor deve ser mysql.

name

O nome da fonte de dados.

Não

STRING

Nenhum

Nenhuma.

hostname

O endereço IP ou nome do host do banco de dados MySQL.

Sim

STRING

Nenhum

Recomendamos especificar um endereço VPC.

Nota

Se o banco de dados MySQL e o Realtime Compute for Apache Flink não estiverem na mesma VPC, estabeleça uma conexão de rede entre VPCs ou use um endpoint público para acessar o banco de dados. Para mais informações, consulte Manage and operate workspaces e How can a fully managed Flink cluster access the Internet?.

username

O nome de usuário do service de banco de dados MySQL.

Sim

STRING

Nenhum

Nenhuma.

password

A senha do service de banco de dados MySQL.

Sim

STRING

Nenhum

Nenhuma.

tables

As tabelas de dados MySQL a serem sincronizadas.

Sim

STRING

Nenhum

  • Este parâmetro aceita expressões regulares para ler dados de várias tabelas.

  • Use vírgulas para separar múltiplas expressões regulares.

Nota
  • Não utilize os caracteres de correspondência de início de string ^ e fim de string $ na expressão regular. No VVR 11.2, o ponto divide a expressão regular para obter a parte do banco de dados. Os caracteres de início e fim tornarão a expressão regular do banco de dados resultante inutilizável. Por exemplo, altere ^db.user_[0-9]+$ para db.user_[0-9]+.

  • O ponto separa o nome do banco de dados e o nome da tabela. Para usar um ponto como coringa, escape-o com uma barra invertida. Exemplo: db0.\., db1.user_table_[0-9]+, db[1-2].[app|web]order_\..

tables.exclude

As tabelas a serem excluídas da sincronização.

Não

STRING

Nenhum

  • Este parâmetro aceita expressões regulares para excluir múltiplas tabelas.

  • Use vírgulas para separar múltiplas expressões regulares.

Nota

O ponto separa o nome do banco de dados e o nome da tabela. Para usar um ponto como coringa, escape-o com uma barra invertida. Exemplo: db0.\., db1.user_table_[0-9]+, db[1-2].[app|web]order_\..

port

O número da porta do service de banco de dados MySQL.

Não

INTEGER

3306

Nenhuma.

schema-change.enabled

Define se eventos de alteração de esquema devem ser enviados.

Não

BOOLEAN

true

Nenhuma.

server-id

Um ID numérico ou intervalo para o cliente de banco de dados usado na sincronização.

Não

STRING

Um valor aleatório entre 5400 e 6400 é gerado.

Este ID deve ser globalmente único dentro do cluster MySQL. Defina um ID diferente para cada job que se conecta ao mesmo banco de dados. Este parâmetro também aceita um formato de intervalo de ID, como 5400-5408.

Nota

Quando a leitura incremental está ativada, há suporte para leitura concorrente. Nesse caso, recomendamos definir um intervalo de IDs para que cada leitor concorrente utilize um ID diferente.

jdbc.properties.*

Parâmetros de conexão personalizados na URL JDBC.

Não

STRING

Nenhum

É possível passar parâmetros de conexão personalizados. Por exemplo, para não usar o protocolo SSL, configure 'jdbc.properties.useSSL' = 'false'.

Para mais informações sobre os parâmetros de conexão suportados, consulte MySQL Configuration Properties.

debezium.*

Parâmetros personalizados do Debezium para leitura de logs binários.

Não

STRING

Nenhum

É possível passar parâmetros personalizados do Debezium. Por exemplo, use 'debezium.event.deserialization.failure.handling.mode'='ignore' para especificar a lógica de tratamento de erros de análise.

Aviso

Não modifique os parâmetros do Debezium arbitrariamente. Isso pode fazer com que o conector leia dados incorretamente. Por exemplo, não é permitido configurar o parâmetro debezium.binlog.buffer.size.

scan.incremental.snapshot.chunk.size

O tamanho de cada chunk em número de linhas.

Não

INTEGER

8096

As tabelas MySQL são divididas em vários chunks para leitura. Os dados de um chunk ficam em cache na memória antes de serem totalmente lidos.

Quanto menos linhas cada chunk contiver, maior será o número total de chunks na tabela. Embora isso reduza a granularidade da recuperação de falhas, pode levar a erros OOM e diminuir o throughput geral. Portanto, avalie o compromisso e defina um tamanho de chunk razoável.

scan.snapshot.fetch.size

O número máximo de registros a serem buscados por vez durante a leitura completa dos dados de uma tabela.

Não

INTEGER

1024

Nenhuma.

scan.startup.mode

O modo de inicialização para consumo de dados.

Não

STRING

initial

Valores válidos:

  • initial (padrão): Na primeira inicialização ou em uma inicialização sem estado, o conector varre todos os dados históricos e depois lê os dados mais recentes do log binário.

  • latest-offset: Na primeira inicialização ou em uma inicialização sem estado, o conector não varre os dados históricos. Ele começa a ler a partir do final do log binário, ou seja, lê apenas as alterações mais recentes feitas após o início do conector.

  • earliest-offset: O conector não varre os dados históricos. Ele começa a ler a partir do log binário disponível mais antigo.

  • specific-offset: O conector não varre os dados históricos. Ele inicia a partir de um offset específico do log binário. Especifique o offset configurando tanto scan.startup.specific-offset.file quanto scan.startup.specific-offset.pos, ou configurando apenas scan.startup.specific-offset.gtid-set para iniciar a partir de um conjunto GTID específico.

  • timestamp: O conector não varre os dados históricos. Ele começa a ler o log binário a partir de um carimbo de data/hora especificado. O carimbo é definido por scan.startup.timestamp-millis em milissegundos.

Importante

Para os modos de inicialização earliest-offset, specific-offset e timestamp, se o esquema da tabela no momento da inicialização for diferente do esquema no momento do offset inicial especificado, o job reportará um erro devido à incompatibilidade de esquemas. Em outras palavras, ao usar esses três modos de inicialização, garanta que o esquema da tabela correspondente não sofra alterações entre a posição de consumo do log binário especificada e o momento de início do job.

scan.startup.specific-offset.file

O nome do arquivo de log binário para o offset inicial ao usar o modo de inicialização specific-offset.

Não

STRING

Nenhum

Ao usar este parâmetro, defina scan.startup.mode como specific-offset. Exemplo de formato de nome de arquivo: mysql-bin.000003.

scan.startup.specific-offset.pos

O offset dentro do arquivo de log binário especificado para o offset inicial ao usar o modo de inicialização specific-offset.

Não

INTEGER

Nenhum

Ao usar este parâmetro, defina scan.startup.mode como specific-offset.

scan.startup.specific-offset.gtid-set

O conjunto GTID para o offset inicial ao usar o modo de inicialização specific-offset.

Não

STRING

Nenhum

Ao usar este parâmetro, defina scan.startup.mode como specific-offset. Exemplo de formato de conjunto GTID: 24DA167-0C0C-11E8-8442-00059A3C7B00:1-19.

scan.startup.timestamp-millis

O carimbo de data/hora em milissegundos para o offset inicial ao usar o modo de inicialização timestamp.

Não

LONG

Nenhum

Ao usar este parâmetro, defina scan.startup.mode como timestamp. A unidade do carimbo de data/hora é milissegundos.

Importante

Ao especificar um horário, o MySQL CDC tenta ler o evento inicial de cada arquivo de log binário para determinar seu carimbo de data/hora. Em seguida, localiza o arquivo de log binário correspondente ao horário especificado. Certifique-se de que o arquivo de log binário correspondente ao carimbo especificado não tenha sido limpo do banco de dados e possa ser lido.

server-time-zone

O fuso horário da sessão usado pelo banco de dados.

Não

STRING

Se você não especificar este parâmetro, o sistema usará o fuso horário do ambiente de execução do job Flink como fuso horário do servidor de banco de dados. Este é o fuso horário da zona selecionada.

Exemplo: Asia/Shanghai. Este parâmetro controla como o tipo TIMESTAMP no MySQL é convertido para o tipo STRING. Para mais informações, consulte Debezium temporal values.

scan.startup.specific-offset.skip-events

O número de eventos de log binário a serem ignorados ao ler a partir de um offset especificado.

Não

INTEGER

Nenhum

Ao usar este parâmetro, defina scan.startup.mode como specific-offset.

scan.startup.specific-offset.skip-rows

O número de alterações de linha a serem ignoradas ao ler a partir de um offset especificado. Um único evento de log binário pode corresponder a várias alterações de linha.

Não

INTEGER

Nenhum

Ao usar este parâmetro, defina scan.startup.mode como specific-offset.

connect.timeout

O tempo máximo de espera para que uma conexão com o servidor de banco de dados MySQL atinja o timeout antes de tentar novamente.

Não

DURATION

30 s

Nenhuma.

connect.max-retries

O número máximo de tentativas após uma falha de conexão com o service de banco de dados MySQL.

Não

INTEGER

3

Nenhuma.

connection.pool.size

O tamanho do pool de conexões do banco de dados.

Não

INTEGER

20

O pool de conexões do banco de dados serve para reutilizar conexões, o que reduz o número total de conexões com o banco.

heartbeat.interval

O intervalo no qual a fonte avança o offset do log binário usando eventos de heartbeat.

Não

DURATION

30s

Eventos de heartbeat servem para avançar o offset do log binário na fonte. Isso é muito útil para tabelas no MySQL que são atualizadas com pouca frequência. Para tais tabelas, o offset do log binário não consegue avançar automaticamente. Eventos de heartbeat podem empurrar o offset do log binário para frente, evitando problemas causados por um offset expirado. Um offset de log binário expirado pode fazer com que o job falhe e se torne irrecuperável, exigindo uma reinicialização sem estado.

rds.region-id

O ID da região da instância Alibaba Cloud ApsaraDB RDS for MySQL.

Obrigatório ao usar o recurso de leitura de logs arquivados do OSS.

STRING

Nenhum

Para mais informações sobre IDs de região, consulte Regions and zones.

Importante

Como a string GTID para MySQL CDC é gerada aleatoriamente e não aumenta monotonicamente como os offsets de arquivos de log binário, localizar um GTID em um arquivo exige baixar e analisar todos os logs arquivados do OSS. Esse processo consome muitos recursos e tempo, tornando inviáveis os recursos que dependem de offsets GTID. Portanto, o recurso de logs arquivados do OSS suporta apenas a inicialização a partir de um carimbo de data/hora especificado ou de um offset de arquivo de log binário especificado. Ele não suporta inicialização a partir de um GTID especificado, nem cenários com failovers primário/secundário nos logs arquivados, pois os failovers do MySQL dependem de GTIDs. Avalie cuidadosamente este recurso antes de usá-lo.

rds.access-key-id

O AccessKey ID da conta Alibaba Cloud ApsaraDB RDS for MySQL.

Obrigatório ao usar o recurso de leitura de logs arquivados do OSS.

STRING

Nenhum

Para mais informações, consulte How do I view the AccessKey ID and AccessKey secret?

Importante

Para evitar que suas informações de AccessKey vazem, use o recurso de gerenciamento de segredos para especificar o AccessKey ID. Para mais informações, consulte Manage variables.

rds.access-key-secret

O AccessKey secret da conta Alibaba Cloud ApsaraDB RDS for MySQL.

Obrigatório ao usar o recurso de leitura de logs arquivados do OSS.

STRING

Nenhum

Para mais informações, consulte How do I view the AccessKey ID and AccessKey secret?

Importante

Para evitar que suas informações de AccessKey vazem, use o recurso de gerenciamento de segredos para especificar o AccessKey secret. Para mais informações, consulte Manage variables.

rds.db-instance-id

O ID da instância Alibaba Cloud ApsaraDB RDS for MySQL.

Obrigatório ao usar o recurso de leitura de logs arquivados do OSS.

STRING

Nenhum

Nenhuma.

rds.main-db-id

O número do banco de dados primário da instância Alibaba Cloud ApsaraDB RDS for MySQL.

Não

STRING

Nenhum

Para mais informações sobre como obter o número do banco de dados primário, consulte ApsaraDB RDS for MySQL log backup.

Nota

Se este parâmetro não for especificado, o VVR 11.7 e versões posteriores consultarão automaticamente o número do banco de dados primário com base nas informações de conexão do ApsaraDB RDS for MySQL.

rds.download.timeout

O tempo limite para baixe um único log arquivado do OSS.

Não

DURATION

60s

Nenhuma.

rds.endpoint

O endpoint de service para obter informações de log binário do OSS.

Não

STRING

Nenhum

Para mais informações sobre os valores válidos, consulte Endpoints.

rds.binlog-directory-prefix

O prefixo do diretório para armazenar arquivos de log binário.

Não

STRING

rds-binlog-

Nenhuma.

rds.use-intranet-link

Define se a rede interna deve ser usada para baixe arquivos de log binário.

Não

BOOLEAN

true

Nenhuma.

rds.binlog-directories-parent-path

O caminho absoluto do diretório pai para armazenar arquivos de log binário.

Não

STRING

Nenhum

Nenhuma.

chunk-meta.group.size

O tamanho dos metadados do chunk.

Não

INTEGER

1000

Se os metadados forem maiores que este valor, eles serão divididos em várias partes para transmissão.

chunk-key.even-distribution.factor.lower-bound

O limite inferior do fator de distribuição de chunks para fragmentação uniforme.

Não

DOUBLE

0,05

Se o fator de distribuição for menor que este valor, a fragmentação não uniforme será usada.

Fator de distribuição de chunks = (MAX(chunk-key) - MIN(chunk-key) + 1) / Número total de linhas de dados.

chunk-key.even-distribution.factor.upper-bound

O limite superior do fator de distribuição de chunks para fragmentação uniforme.

Não

DOUBLE

1.000,0

Se o fator de distribuição for maior que este valor, a fragmentação não uniforme será usada.

Fator de distribuição de chunks = (MAX(chunk-key) - MIN(chunk-key) + 1) / Número total de linhas de dados.

scan.incremental.close-idle-reader.enabled

Define se leitores ociosos devem ser fechados após a conclusão do snapshot.

Não

BOOLEAN

false

Para que esta configuração tenha efeito, defina execution.checkpointing.checkpoints-after-tasks-finish.enabled como true.

scan.only.deserialize.captured.tables.changelog.enabled

Na fase incremental, define se apenas os eventos de alteração das tabelas especificadas devem ser desserializados.

Não

BOOLEAN

  • O valor padrão é false nas versões VVR 8.x.

  • O valor padrão é true no VVR 11.1 e posteriores.

Valores válidos:

  • true: Desserializa apenas os dados de alteração das tabelas alvo para acelerar a leitura do log binário.

  • false (padrão): Desserializa os dados de alteração de todas as tabelas.

scan.parallel-deserialize-changelog.enabled

Na fase incremental, define se múltiplas threads devem ser usadas para analisar eventos de alteração.

Não

BOOLEAN

false

Valores válidos:

  • true: Usa múltiplas threads na fase de desserialização de eventos de alteração, mantendo a ordem dos eventos de log binário para acelerar a leitura.

  • false (padrão): Usa uma única thread na fase de desserialização de eventos.

Nota

Suportado apenas no VVR 8.0.11 e posteriores.

scan.parallel-deserialize-changelog.handler.size

O número de manipuladores de eventos ao usar múltiplas threads para analisar eventos de alteração.

Não

INTEGER

2

Nota

Suportado apenas no VVR 8.0.11 e posteriores.

metadata-column.include-list

As colunas de metadados a serem passadas para o downstream.

Não

STRING

Nenhum

Os metadados disponíveis incluem op_ts, es_ts, query_log, file e pos. Use vírgulas para separar múltiplas colunas de metadados.

Nota

O conector YAML MySQL CDC não exige nem suporta a adição de colunas de metadados de nome do banco de dados, nome da tabela e op_type. Use diretamente __data_event_type__ em uma expressão Transform para obter o tipo de dado de alteração, ou use __schema_name__ e __table_name__ para obter o nome do banco de dados e da tabela.

Importante
  • A coluna de metadados file representa o arquivo de log binário onde os dados estão localizados. Ela é "" durante a fase completa e o nome do arquivo de log binário durante a fase incremental. A coluna de metadados pos representa o offset dos dados no arquivo de log binário. Ela é "0" durante a fase completa e o offset dos dados no arquivo de log binário durante a fase incremental. Essas duas colunas de metadados são suportadas a partir do VVR 11.5.

  • A coluna de metadados es_ts representa a hora de início da transação correspondente para o changelog no MySQL. Ela é suportada apenas para MySQL 8.0.x. Não adicione esta coluna de metadados ao usar versões anteriores do MySQL.

  • O carimbo de data/hora op_ts tem precisão de segundos, enquanto o carimbo es_ts tem precisão de milissegundos.

scan.newly-added-table.enabled

Ao reiniciar a partir de um checkpoint, define se tabelas recém-adicionadas que não foram correspondidas na inicialização anterior devem ser sincronizadas, ou se tabelas que não correspondem mais devem ser removidas do estado.

Não

BOOLEAN

false

Isso entra em vigor ao reiniciar a partir de um checkpoint ou savepoint.

Importante

Durante a fase de leitura completa de dados, não é possível salve um savepoint, adicionar ou remover uma tabela da tabela de origem e reiniciar o job a partir desse savepoint. Essa ação causa falha na leitura dos dados.

scan.binlog.newly-added-table.enabled

Na fase incremental, define se dados de tabelas recém-adicionadas que forem correspondidas devem ser enviados.

Não

BOOLEAN

false

Não pode ser habilitado simultaneamente com scan.newly-added-table.enabled.

scan.incremental.snapshot.chunk.key-column

Especifica uma coluna de determinadas tabelas para ser usada como coluna de divisão para fragmentação durante a fase de snapshot.

Não

STRING

Nenhum

  • Use dois pontos : para conectar o nome da tabela e o nome da coluna e definir uma regra. O nome da tabela pode ser uma expressão regular. É possível definir várias regras separando-as com ponto e vírgula ;. Exemplo: db1.user_table_[0-9]+:col1;db[1-2].[app|web]_order_\\.*:col2.

  • Obrigatório para tabelas sem chave primária. A coluna selecionada deve ser de um tipo não nulo (NOT NULL). Opcional para tabelas com chave primária. Apenas uma coluna pode ser selecionada da chave primária.

scan.parse.online.schema.changes.enabled

Na fase incremental, define se deve haver tentativa de analisar eventos DDL de alteração sem bloqueio do RDS.

Não

BOOLEAN

false

Valores válidos:

  • true: Analisa eventos DDL de alteração sem bloqueio do RDS.

  • false (padrão): Não analisa eventos DDL de alteração sem bloqueio do RDS.

Este é um recurso experimental. Antes de executar uma alteração online sem bloqueio, faça um snapshot do job Flink para recuperação.

Nota

Suportado apenas no VVR 11.0 e posteriores.

scan.incremental.snapshot.backfill.skip

Define se o backfill deve ser ignorado durante a fase de leitura de snapshot.

Não

BOOLEAN

false

Valores válidos:

  • true: Ignora o backfill durante a fase de leitura de snapshot.

  • false (padrão): Não ignora o backfill durante a fase de leitura de snapshot.

O backfill aplica-se apenas durante a consulta de snapshot de um único chunk e não cobre toda a fase de leitura completa. Quando o backfill é ignorado, a consulta de snapshot de cada chunk lê os dados mais recentes da tabela naquele instante; atualizações que ocorrem em um chunk após sua leitura não são mescladas durante a fase de leitura completa e são lidas do Binlog após entrar na fase incremental. Por exemplo, uma atualização no chunk5 que ocorre enquanto o chunk5 está sendo capturado no snapshot é refletida diretamente no snapshot do chunk5; se o chunk5 for atualizado depois que o leitor avançou para o chunk80, a atualização será aplicada posteriormente a partir do Binlog durante a fase incremental.

Importante

Quando habilitado, as alterações que ocorrem durante ou após a varredura de um chunk ainda são entregues pelo Binlog na fase incremental e podem ser duplicadas. Apenas a semântica at-least-once é garantida. Habilite isso apenas quando o sink downstream suportar gravações idempotentes por chave primária.

Nota

Suportado apenas no VVR 11.1 e posteriores.

treat-tinyint1-as-boolean.enabled

Define se o tipo TINYINT(1) deve ser tratado como um tipo Boolean.

Não

BOOLEAN

true

Valores válidos:

  • true (padrão): Trata o tipo TINYINT(1) como um tipo Boolean.

  • false: Não trata o tipo TINYINT(1) como um tipo Boolean.

treat-timestamp-as-datetime-enabled

Define se o tipo TIMESTAMP deve ser tratado como um tipo DATETIME.

Não

BOOLEAN

false

Valores válidos:

  • true: Trata o tipo TIMESTAMP do MySQL como um tipo DATETIME e o mapeia para o tipo TIMESTAMP do CDC.

  • false (padrão): Mapeia o tipo TIMESTAMP do MySQL para o tipo TIMESTAMP_LTZ do CDC.

O tipo TIMESTAMP do MySQL armazena o horário UTC e é afetado pelo fuso horário. O tipo DATETIME do MySQL armazena o horário literal e não é afetado pelo fuso horário.

Quando habilitado, este parâmetro converte dados do tipo TIMESTAMP do MySQL para o tipo DATETIME com base no server-time-zone.

include-comments.enabled

Define se comentários de tabelas e colunas devem ser sincronizados.

Não

BOOELEAN

false

Valores válidos:

  • true: Sincroniza comentários de tabelas e colunas.

  • false (padrão): Não sincroniza comentários de tabelas e colunas.

Habilitar esta opção aumenta o uso de memória do job.

scan.incremental.snapshot.unbounded-chunk-first.enabled

Define se chunks ilimitados devem ser despachados primeiro durante a fase de leitura de snapshot.

Não

BOOELEAN

false

Valores válidos:

  • true: Despacha chunks ilimitados primeiro durante a fase de leitura de snapshot.

  • false (padrão): Não despacha chunks ilimitados primeiro durante a fase de leitura de snapshot.

Este é um recurso experimental. Habilitá-lo pode reduzir o risco de erros OOM no TaskManager ao sincronizar o último chunk durante a fase de snapshot. Adicione este parâmetro antes da primeira inicialização do job.

Nota

Suportado apenas no VVR 11.1 e posteriores.

binlog.session.network.timeout

O timeout de rede para a conexão do log binário.

Não

DURATION

10m

Se definido como 0s, o timeout padrão do servidor MySQL será usado.

Nota

Suportado apenas no VVR 11.5 e posteriores.

scan.rate-limit.records-per-second

Limita o número máximo de registros enviados pela fonte por segundo.

Não

LONG

Nenhum

Aplicável a cenários onde a leitura de dados precisa ser limitada. Esse limite é efetivo tanto nas fases completa quanto incremental.

A métrica numRecordsOutPerSecond da fonte reflete o número de registros produzidos por todo o fluxo de dados por segundo. Ajuste este parâmetro com base nessa métrica.

Na fase de leitura completa de dados, geralmente é necessário reduzir o número de linhas lidas em cada lote. Reduza o valor do parâmetro scan.incremental.snapshot.chunk.size.

Nota

Suportado apenas no VVR 11.5 e posteriores.

include-binlog-meta.enable

Define se as informações originais do log binário do MySQL, como GTID e offset do log binário, devem ser incluídas na mensagem.

Não

Boolean

false

Aplicável a cenários de sincronização original de log binário, como substituir um link de sincronização Canal existente.

Nota

Suportado apenas no VVR 11.6 e posteriores.

scan.binlog.tolerate.gtid-holes

Habilitar este parâmetro ignora lacunas na sequência GTID, permitindo que o job contorne eventos descontínuos e continue executando.

Não

Boolean

false

Antes de habilitar este parâmetro, garanta que o offset inicial do job não tenha expirado. Se o job iniciar a partir de um offset GTID limpo ou expirado, o mecanismo ignorará silenciosamente os logs ausentes, o que levará à perda de dados.

Nota

Este parâmetro é suportado apenas no VVR 11.6 e posteriores.

scan.emit.create-table-events.in-batch.enabled

Define se esquemas de tabela devem ser enviados em lote durante a fase de inicialização do job.

Não

Boolean

false

Este é um recurso experimental. Habilite esta opção quando um único job sincronizar muitas tabelas.

Nota

Este parâmetro é suportado apenas no VVR 11.4 e posteriores.

Reutilizar um catálogo existente

A partir do VVR 11.5, é possível referenciar diretamente um catálogo MySQL integrado criado na página Data Management em um job de ingestão de dados Flink CDC. Isso reduz o esforço manual de escrever propriedades de conexão.

source:
  type: mysql
  using.built-in-catalog: mysql_rds_catalog

Atualmente, os jobs de ingestão de dados suportam a reutilização automática dos seguintes parâmetros de catálogo MySQL:

  • hostname

  • port

  • username

  • password

  • catalog.table.metadata-columns

  • catalog.table.treat-tinyint1-as-boolean

Para substituir qualquer um desses parâmetros reutilizados automaticamente, escreva explicitamente o parâmetro YAML correspondente. O parâmetro escrito explicitamente tem prioridade mais alta.

Mapeamento de tipos

A tabela a seguir mostra o mapeamento de tipos de dados para ingestão de dados.

Tipo de campo MySQL CDC

Tipo de campo CDC

TINYINT(n)

TINYINT

SMALLINT

SMALLINT

TINYINT UNSIGNED

TINYINT UNSIGNED ZEROFILL

YEAR

INT

INT

MEDIUMINT

MEDIUMINT UNSIGNED

MEDIUMINT UNSIGNED ZEROFILL

SMALLINT UNSIGNED

SMALLINT UNSIGNED ZEROFILL

BIGINT

BIGINT

INT UNSIGNED

INT UNSIGNED ZEROFILL

BIGINT UNSIGNED

DECIMAL(20, 0)

BIGINT UNSIGNED ZEROFILL

SERIAL

FLOAT [UNSIGNED] [ZEROFILL]

FLOAT

DOUBLE [UNSIGNED] [ZEROFILL]

DOUBLE

DOUBLE PRECISION [UNSIGNED] [ZEROFILL]

REAL [UNSIGNED] [ZEROFILL]

NUMERIC(p, s) [UNSIGNED] [ZEROFILL] e p <= 38

DECIMAL(p, s)

DECIMAL(p, s) [UNSIGNED] [ZEROFILL] e p <= 38

FIXED(p, s) [UNSIGNED] [ZEROFILL] e p <= 38

BOOLEAN

BOOLEAN

BIT(1)

TINYINT(1)

DATE

DATE

TIME [(p)]

TIME [(p)]

DATETIME [(p)]

TIMESTAMP [(p)]

TIMESTAMP [(p)]

O mapeamento depende do valor do parâmetro treat-timestamp-as-datetime-enabled:

true:TIMESTAMP[(p)]

false:TIMESTAMP_LTZ[(p)]

CHAR(n)

CHAR(n)

VARCHAR(n)

VARCHAR(n)

BIT(n)

BINARY(⌈(n + 7) / 8⌉)

BINARY(n)

BINARY(n)

VARBINARY(N)

VARBINARY(N)

NUMERIC(p, s) [UNSIGNED] [ZEROFILL] e 38 < p <= 65

STRING

Nota

No MySQL, o tipo de dado decimal tem precisão de até 65, mas no Flink a precisão é limitada a 38. Portanto, se você definir uma coluna decimal com precisão maior que 38, mapeie-a para string para evitar perda de precisão.

DECIMAL(p, s) [UNSIGNED] [ZEROFILL] e 38 < p <= 65

FIXED(p, s) [UNSIGNED] [ZEROFILL] e 38 < p <= 65

TINYTEXT

STRING

TEXT

MEDIUMTEXT

LONGTEXT

ENUM

JSON

STRING

Nota

O tipo de dado JSON é convertido em uma string formatada em JSON no Flink.

GEOMETRY

STRING

Nota

Os tipos de dados espaciais no MySQL são convertidos em strings com um formato JSON fixo. Para mais informações, consulte Spatial Data Type Mapping do MySQL.

POINT

LINESTRING

POLYGON

MULTIPOINT

MULTILINESTRING

MULTIPOLYGON

GEOMETRYCOLLECTION

TINYBLOB

BYTES

Nota

Para o tipo de dado BLOB no MySQL, apenas blobs com comprimento não superior a 2.147.483.647 (2**31-1) são suportados.

BLOB

MEDIUMBLOB

LONGBLOB

Definir um server ID para evitar conflitos de consumo de log binário

Quando um job de ingestão de dados lê logs binários, a fonte se registra no MySQL como um cliente de replicação usando um server-id. Se múltiplos jobs ou outros clientes de replicação usarem o mesmo server ID, ocorrerão conflitos de consumo de log binário e os jobs falharão. Observe o seguinte ao configure server IDs:

  • Por padrão, um server-id é um valor único aleatório entre 5400 e 6400. Se múltiplos jobs usarem valores padrão, poderão ocorrer conflitos. Recomendamos configure explicitamente IDs não sobrepostos para os jobs.

  • Se o paralelismo da fonte for maior que 1, configure um intervalo de server ID. O número de IDs disponíveis no intervalo não deve ser menor que o paralelismo, e cada leitor paralelo usa um ID diferente.

Exemplo: Se o paralelismo da fonte for 4, configure um intervalo que contenha quatro IDs.

source:
  type: mysql
  name: MySQL Source
  hostname: <hostname>
  port: 3306
  username: <username>
  password: <password>
  tables: app_db.\.*
  server-id: 5400-5403

sink:
  type: hologres

Se múltiplos jobs lerem da mesma instância MySQL, atribua um intervalo não sobreposto a cada job. Por exemplo, o job A usa 5400-5403, e o job B usa 5404-5407.

Acelerar a leitura de log binário

Ao usar o conector MySQL como fonte de dados de ingestão, ele analisa arquivos de log binário para gerar várias mensagens de alteração durante a fase incremental. Os arquivos de log binário registram todas as alterações de tabela em formato binário. É possível acelerar a análise dos arquivos de log binário das seguintes maneiras.

  • Habilite a análise paralela e filtros de análise (Este recurso requer Realtime Compute for Apache Flink com Ververica Runtime (VVR) 8.0.7 ou posterior. Não está disponível na edição community do conector MySQL CDC.)

    • Ative a opção scan.only.deserialize.captured.tables.changelog.enabled para analisar eventos de alteração apenas para tabelas especificadas.

    • Ative a opção scan.parallel-deserialize-changelog.enabled para usar múltiplas threads na análise do arquivo de log binário e entregar eventos à fila do consumidor em ordem. Ao ativar esta opção, geralmente é necessário aumentar também a CPU do TaskManager.

  • Otimize os parâmetros do Debezium

    debezium.max.queue.size: 162580
    debezium.max.batch.size: 40960
    debezium.poll.interval.ms: 50
    • debezium.max.queue.size: O número máximo de registros que a fila de bloqueio pode conter. Quando o Debezium lê um fluxo de eventos do banco de dados, ele coloca os eventos em uma fila de bloqueio antes de gravá-los no downstream. O valor padrão é 8192.

    • debezium.max.batch.size: O número máximo de eventos que o conector processa em cada iteração. O valor padrão é 2048.

    • debezium.poll.interval.ms: O número de milissegundos que o conector deve esperar antes de solicitar novos eventos de alteração. O valor padrão é 1000 milissegundos, ou 1 segundo.

Exemplo de uso:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: ${mysql.source.table}
  server-id: 7601-7604
  # Debezium configuration
  debezium.max.queue.size: 162580
  debezium.max.batch.size: 40960
  debezium.poll.interval.ms: 50
  # Enable parsing filter
  scan.only.deserialize.captured.tables.changelog.enabled: true

A capacidade de consumo de log binário da Enterprise Edition do MySQL CDC é de 85 MB/s, cerca do dobro da versão open source da comunidade. Quando a velocidade de geração de arquivos de log binário excede 85 MB/s (ou seja, um arquivo de 512 MB a cada 6 segundos), a latência do job Flink continua aumentando. A latência de processamento diminui gradualmente após a desaceleração da velocidade de geração do arquivo de log binário. Se um arquivo de log binário contiver uma transação grande, a latência de processamento pode aumentar temporariamente. A latência diminui após a leitura do log dessa transação.

Diagnosticar latência de dados para otimizar o throughput do job

Se você enfrentar latência de dados durante a fase incremental, analise o problema seguindo estas etapas:

  1. Verifique as métricas currentFetchEventTimeLag e currentEmitEventTimeLag na página Overview. A métrica currentFetchEventTimeLag representa a latência na leitura de dados do log binário. A métrica currentEmitEventTimeLag representa a latência na leitura de dados das tabelas relevantes para o job a partir do log binário.

    Cenário

    Descrição

    currentFetchEventTimeLag está baixo, enquanto currentEmitEventTimeLag está alto e raramente se atualiza.

    Um currentFetchEventTimeLag baixo indica que a extração do log binário do banco de dados é eficiente. No entanto, o log binário contém poucos dados para as tabelas que o job precisa ler. Portanto, currentEmitEventTimeLag raramente se atualiza. Este é um comportamento esperado.

    Tanto currentFetchEventTimeLag quanto currentEmitEventTimeLag estão altos.

    Isso indica que a tabela de origem tem baixo desempenho de leitura. Prossiga para as etapas subsequentes nesta seção para otimização.

  2. A contrapressão pode reduzir a taxa na qual a fonte envia dados para operadores downstream. Você pode observar que sourceIdleTime aumenta periodicamente, e tanto currentFetchEventTimeLag quanto currentEmitEventTimeLag crescem continuamente. Para resolver isso, aumente o paralelismo do nó onde a contrapressão se origina.

  3. Verifique a métrica TM CPU Usage na página CPU e a métrica TM GC Time na página JVM para determinar se há recursos insuficientes de CPU ou memória. Aumente os recursos do job para otimizar o desempenho de leitura.

Ler binlogs arquivados do OSS

Ao usar uma instância ApsaraDB RDS for MySQL como fonte de dados, é possível ler backups de log armazenados no OSS. Se o arquivo correspondente ao carimbo de data/hora ou posição de log binário especificado estiver armazenado no OSS, o Flink extrairá automaticamente o arquivo de log do OSS para o cluster. Se o arquivo estiver armazenado localmente no banco de dados, o Flink mudará automaticamente para a leitura através de uma conexão de banco de dados. Este recurso está disponível apenas no Realtime Compute for Apache Flink e não é suportado na edição community do conector MySQL CDC.

Para habilitar a leitura de backups de log do OSS, configure os parâmetros de conexão do ApsaraDB RDS for MySQL. Exemplo:

source:
  type: mysql
  hostname: <yourHostname>
  port: 3306
  username: <yourUsername>
  password: <yourPassword>
  tables: <yourTables>
  # RDS connection parameters for reading archived binlogs from OSS
  rds.region-id: cn-beijing
  rds.access-key-id: your_access_key_id
  rds.access-key-secret: your_access_key_secret
  rds.db-instance-id: rm-xxxxxxxx # Database instance ID.
  rds.main-db-id: 12345678 # Primary database ID.
  rds.endpoint: rds.aliyuncs.com

Ingestão de dados para sincronização de banco de dados e esquema

Para jobs que envolvem apenas lógica de sincronização de dados, recomendamos executá-los como jobs de ingestão de dados. Jobs de ingestão são profundamente otimizados para cenários de integração de dados. Para instruções de uso, consulte Get Started with Flink CDC Data Ingestion e Develop a Flink CDC Data Ingestion Job (Public Preview).

O código a seguir mostra um exemplo de job de ingestão de dados que sincroniza todo o banco de dados app_db do MySQL para o Hologres, incluindo alterações de esquema subsequentes do banco de dados upstream:

source:
  type: mysql
  hostname: <hostname>
  port: 3306
  username: ${secret_values.mysqlusername}
  password: ${secret_values.mysqlpassword}
  tables: app_db.\.*
  server-id: 5400-5404

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

Descoberta de novas tabelas na ingestão de dados

O conector de ingestão de dados MySQL fornece opções de configuração para suportar a descoberta de novas tabelas em dois cenários diferentes.

Parâmetro

Descrição

Observações

scan.newly-added-table.enabled

Quando um job reinicia a partir de um checkpoint, esta opção sincroniza tabelas que não foram descobertas durante a inicialização anterior. Ela lê dados de snapshot e incrementais dessas novas tabelas.

Esta opção é suportada apenas quando scan.startup.mode está definido como initial. Não tem efeito em outros modos de inicialização.

scan.binlog.newly-added-table.enabled

Durante a fase incremental, esta opção sincroniza automaticamente dados de tabelas recém-descobertas.

  • Recomendamos habilitar esta opção ao iniciar o job pela primeira vez. O job analisará automaticamente instruções DDL CREATE TABLE e sincronizará os dados downstream. Se você habilitar esta opção e reiniciar o job após a tabela do banco de dados já ter sido criada, isso pode levar a dados incompletos.

  • No modo de inicialização initial, operações DDL não são sincronizadas downstream até que a fase de snapshot seja concluída. Tabelas criadas durante a fase de snapshot não podem ser sincronizadas automaticamente mesmo se scan.binlog.newly-added-table.enabled estiver habilitado.

Importante
  • Durante a fase de leitura completa, não há suporte para salve um savepoint e, em seguida, adicionar ou remover tabelas de origem antes de reiniciar a partir do savepoint. Fazer isso impede que o job leia dados corretamente.

  • Não habilite scan.newly-added-table.enabled e scan.binlog.newly-added-table.enabled simultaneamente. Habilitar ambos pode causar duplicação de dados.