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_imagedeve 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_imagedeve 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_imagedeve 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_orderno 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
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 |
Nota
|
tables.exclude | As tabelas a serem excluídas da sincronização. | Não | STRING | Nenhum |
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:
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: |
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: |
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 |
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 |
| Valores válidos:
|
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:
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 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 Importante
|
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.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 |
|
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:
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:
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:
|
treat-timestamp-as-datetime-enabled | Define se o tipo TIMESTAMP deve ser tratado como um tipo DATETIME. | Não | BOOLEAN | false | Valores válidos:
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:
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:
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 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 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
|
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.enabledpara analisar eventos de alteração apenas para tabelas especificadas.Ative a opção
scan.parallel-deserialize-changelog.enabledpara 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 aCPU do TaskManager.
-
Otimize os parâmetros do Debezium
debezium.max.queue.size: 162580 debezium.max.batch.size: 40960 debezium.poll.interval.ms: 50debezium.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:
-
Verifique as métricas
currentFetchEventTimeLagecurrentEmitEventTimeLagna página Overview. A métricacurrentFetchEventTimeLagrepresenta a latência na leitura de dados do log binário. A métricacurrentEmitEventTimeLagrepresenta a latência na leitura de dados das tabelas relevantes para o job a partir do log binário.Cenário
Descrição
currentFetchEventTimeLagestá baixo, enquantocurrentEmitEventTimeLagestá alto e raramente se atualiza.Um
currentFetchEventTimeLagbaixo 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,currentEmitEventTimeLagraramente se atualiza. Este é um comportamento esperado.Tanto
currentFetchEventTimeLagquantocurrentEmitEventTimeLagestã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.
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.
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 |
| 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 |
| Durante a fase incremental, esta opção sincroniza automaticamente dados de tabelas recém-descobertas. |
|
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.enabledescan.binlog.newly-added-table.enabledsimultaneamente. Habilitar ambos pode causar duplicação de dados.