Um catálogo do MySQL permite acessar diretamente as tabelas em uma instância do MySQL a partir do console do Realtime Compute for Apache Flink e usá-las em implantações do Flink SQL. Este tópico mostra como criar e usar um catálogo do MySQL.
Contexto
Um catálogo do MySQL oferece os seguintes recursos:
Acessar tabelas em uma instância do MySQL diretamente, sem registro manual por meio de instruções DDL, melhorando a eficiência e a precisão do desenvolvimento.
Utilizar tabelas do catálogo do MySQL como tabelas de origem CDC, tabelas de destino ou tabelas de dimensão em implantações do Flink SQL.
Oferecer suporte a ApsaraDB RDS for MySQL, PolarDB for MySQL e bancos de dados MySQL autogerenciados.
Permitir acesso direto a tabelas lógicas para tabelas com shard.
Integrar-se a jobs de ingestão de dados do Flink CDC para sincronizar o banco de dados completo, fazer a sincronização mesclada de tabelas com shard e sincronizar alterações de esquema em fontes de dados do MySQL.
Limitações
O Realtime Compute for Apache Flink e a instância do MySQL devem estar na mesma VPC. Para conectar VPCs diferentes ou acessar pela internet, estabeleça a conectividade de rede. Para mais informações, consulte network connectivity.
Não modifique a configuração do catálogo após a criação. Para alterá-la, exclua o catálogo e crie-o novamente.
Não é possível criar bancos de dados ou tabelas com o Flink.
-
Quando usadas como origem, essas tabelas suportam apenas leitura de stream, e não leitura em lote.
NotaAntes de usar uma tabela de um catálogo do MySQL como tabela de origem CDC, ative o log binário (Binlog) no ApsaraDB RDS for MySQL, PolarDB for MySQL ou no banco de dados MySQL autogerenciado. Para mais informações, consulte Configure a MySQL database.
-
O catálogo não consegue identificar tabelas que usam sintaxe específica do PolarDB em suas instruções DDL.
Por exemplo,
PARTITION BY KEY(. No Ververica Runtime (VVR) 8.0.7 e versões posteriores, não é possível usar views como tabelas do Flink.
Há suporte apenas para as versões 5.7 e 8.0.x do MySQL.
Crie um catálogo do MySQL
Crie um catálogo do MySQL usando o console ou um comando SQL.
Console (recommended)
-
Acesse a página Data Management.
Faça login no console do Realtime Compute for Apache Flink. Na coluna Actions do workspace a ser gerenciado, clique em Console.
No painel de navegação à esquerda, clique em Data Management.
Clique em Create Catalog, selecione MySQL e, em seguida, clique em Next.
-
Configure os parâmetros.
ImportanteNão modifique esses parâmetros de configuração após a criação. Para fazer alterações, exclua e recrie o catálogo.
Parâmetro
Descrição
Obrigatório
catalogname
Nome do catálogo do MySQL.
Sim
hostname
Endereço IP ou hostname do banco de dados MySQL.
NotaPara conectar VPCs diferentes ou acessar pela internet, estabeleça a conectividade de rede. Para mais informações, consulte network connectivity.
Sim
port
Número da porta do banco de dados MySQL. Padrão: 3306.
Não
default-database
Nome do banco de dados MySQL padrão.
Sim
username
Nome de usuário do banco de dados MySQL.
Sim
password
Senha do banco de dados MySQL.
Para evitar a exposição de segredos em texto simples, recomendamos usar uma variável. O exemplo usa uma variável chamada mysqlpw. Para mais informações, consulte Create a variable.
Sim
-
Clique em OK.
O catálogo criado aparece na área Catalogs, à esquerda.
SQL command
-
Acesse a página Scripts.
Faça login no console do Realtime Compute for Apache Flink. Na coluna Actions do workspace a ser gerenciado, clique em Console.
No painel de navegação à esquerda, clique em .
Clique em
, clique em New Script, insira o File Name e o Storage Location e, em seguida, clique em Save.-
Insira o seguinte código.
CREATE CATALOG YourCatalogName WITH( 'type' = 'mysql', 'hostname' = 'rm-bp1gcn0q0j0******.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'usertest', 'password' = '${secret_values.mysqlpw}', 'default-database' = 'flinktest', 'catalog.table.metadata-columns'='table_name' );Parâmetro
Descrição
Obrigatório
YourCatalogName
Nome do catálogo do MySQL.
Sim
type
Tipo do catálogo. Defina como
mysql.Sim
hostname
Endereço IP ou hostname do banco de dados MySQL.
NotaPara conectar VPCs diferentes ou acessar pela internet, estabeleça a conectividade de rede. Para mais informações, consulte network connectivity.
Sim
port
Número da porta do banco de dados MySQL. Padrão: 3306.
Não
default-database
Nome do banco de dados MySQL padrão.
Sim
username
Nome de usuário do banco de dados MySQL.
Sim
password
Senha do banco de dados MySQL.
Para evitar a exposição de segredos em texto simples, recomendamos usar uma variável. O exemplo usa uma variável chamada mysqlpw. Para mais informações, consulte Create a variable.
Sim
property-version
Versão do esquema de propriedades do catálogo. Defina como
0(padrão) ou1(recomendado).Diferentes versões podem ter suporte a propriedades e valores padrão distintos. Consulte as descrições das propriedades para mais detalhes.
Nota-
Compatível apenas com o VVR 8.0.6 e versões posteriores.
-
No VVR 11.1 e versões posteriores, o valor padrão é 1.
Não
catalog.table.metadata-columns
Especifique as colunas de metadados de uma tabela de origem CDC do MySQL a serem adicionadas ao esquema da tabela durante a consulta. Por padrão, nenhuma coluna de metadados é adicionada.
Separe várias colunas de metadados com ponto e vírgula (;), por exemplo:
op_ts;table_name;database_name.Nota-
Compatível apenas com o VVR 6.0.5 e versões posteriores.
-
Definir esta propriedade adiciona as colunas de metadados especificadas ao esquema. Como essas colunas são específicas de tabelas de origem CDC do MySQL, use tabelas deste catálogo apenas como tabelas de origem, e não como tabelas de destino ou de dimensão.
Não
catalog.table.treat-tinyint1-as-boolean
Ao buscar um esquema de tabela, especifica se os tipos
TinyInt(1)eBooleando MySQL devem ser mapeados paraBooleando Flink. Valores válidos:-
true: mapeia para Boolean. -
false: mapeia para TINYINT.
Valor padrão:
-
Se property-version for
0, o valor padrão étrue. -
Se property-version for
1, o valor padrão éfalse.
Nota-
Compatível apenas com o VVR 8.0.4 e versões posteriores.
-
Não recomendamos usar
TinyInt(1)no MySQL para armazenar valores diferentes de 0 e 1. Escolha um mapeamento de tipo apropriado. Para mais informações, consulte Type mapping.
Não
-
-
Selecione a instrução CREATE CATALOG e clique em Run ao lado do número da linha à esquerda.
A mensagem
The following statement has been executed successfully!indica que o catálogo foi criado.A instrução SQL no editor,
CREATE CATALOG myCatalog, tem parâmetros de configuração que incluemtype=mysql, hostname=rm-bp1gcn0q0j0.mysql.rds.aliyuncs.com,port=3306,username=usertest,password=${secret_values.mysqlpw}(referência a uma variável),default-database=flinktestecatalog.table.metadata-columns=table_name.
Visualize e exclua um catálogo do MySQL
Console (recommended)
Na página Data Management, visualize o Name e o Type dos catálogos criados na Catalog List.
-
Visualize: na coluna Actions do catálogo, clique em View para ver os bancos de dados e as tabelas do catálogo.
Os detalhes do esquema da tabela não exibem comentários de campos.
-
Exclua: na coluna Actions do catálogo, clique em Delete.
Esta operação exclui apenas o catálogo, e não as tabelas subjacentes no service associado. As implantações ativas que usam tabelas do catálogo não são afetadas. No entanto, reimplantar ou reiniciar essas implantações causará um erro, pois as tabelas não serão mais encontradas. Prossiga com cautela.
SQL command
-
No editor da página Scripts, insira os seguintes comandos.
-- View the table schema in Flink. Field comments are not displayed. DESCRIBE `<catalogname>`.`<dbname>`.`<tablename>`; -- Delete the catalog. DROP CATALOG `<catalogname>`;NotaEsta operação exclui apenas o catálogo, e não as tabelas subjacentes no service associado. As implantações em execução que usam tabelas do catálogo não são afetadas. No entanto, se você reimplantar ou reiniciar uma implantação, um erro será relatado, pois a tabela não poderá ser encontrada. Prossiga com cautela.
-
Selecione o comando, clique com o botão direito e escolha Run.
A execução da instrução
DESCRIBEretorna o esquema de uma tabela, incluindo campos comoorderkey,custkey,order_statusetotal_price, juntamente com seus tipos de dados e propriedades.
Usar um catálogo do MySQL
Ler de uma tabela de origem do MySQL
INSERT INTO `<othersinktable>`
SELECT ...
FROM `<mysqlcatalog>`.`<dbname>`.`<tablename>` /*+ OPTIONS('server-id' = '6000-6008') */;
Ao usar uma tabela de um catálogo do MySQL como tabela de origem CDC, recomendamos usar SQL hints para especificar um server-id diferente para cada implantação do Flink SQL. Se a tabela de origem exigir um paralelismo maior, configure o server-id como um intervalo com tamanho maior ou igual ao paralelismo.
Ler de tabelas lógicas para tabelas com shard
Um catálogo do MySQL é compatível com o uso de expressões regulares para nomes de bancos de dados e tabelas, permitindo a leitura de dados de tabelas com shard como uma única tabela lógica.
Por exemplo, um banco de dados com shard contém várias tabelas, como user01, user02 e user99, distribuídas em bancos de dados como db01 a db10. Se todas as tabelas tiverem esquemas compatíveis, use expressões regulares para os nomes do banco de dados e da tabela para acessar todas as tabelas de usuário com shard.
SELECT ... FROM `db.*`.`user.*` /*+ OPTIONS('server-id'='6000-6018') */;
Uma tabela lógica para tabelas com shard retorna dois campos de sistema adicionais: _db_name (STRING) e _table_name (STRING). Esses campos, juntamente com a chave primária original, formam uma nova chave primária composta que garante a unicidade. Por exemplo, se a chave primária das tabelas user01 a user99 for id, a chave primária composta da tabela lógica user será (_db_name, _table_name, id).
O catálogo do MySQL é compatível com o uso de expressões regulares para corresponder a várias tabelas a serem sincronizadas, o que permite a sincronização mesclada de tabelas com shard. Para ver um exemplo, consulte Merge and synchronize sharded tables.
Usar a ingestão de dados do Flink CDC para sincronizar dados e alterações de esquema do MySQL em tempo real
A ingestão de dados do Flink CDC é compatível com a sincronização de tabela única, a sincronização de alterações de esquema, a sincronização mesclada de tabelas com shard e a sincronização com colunas calculadas personalizadas. Também oferece suporte à sincronização em tempo real de esquemas e dados no nível do banco de dados, incluindo alterações de esquema. Para exemplos e detalhes, consulte MySQL YAML connector.
# Single-table sync: synchronizes table-level schema changes and data changes in real time.
source:
type: mysql
using.built-in-catalog: mysql-catalog
tables: "<dbname>.<tablename>"
server-id: "6000-6018"
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: xxx
xxx
route:
- source-table: "<db-name>.<table-name>"
sink-table: "<target-db-name>.<target-table-name>"
# Full-database sync: synchronizes database-level schema changes and data changes in real time.
source:
type: mysql
using.built-in-catalog: mysql-catalog
tables: '<dbname>.\.*'
server-id: "6000-6018"
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: xxx
xxx
route:
- source-table: '<dbname>.\.*'
sink-table: "<target-db-name>.<>"
replace-symbol: "<>"
Por exemplo, para sincronizar dados do MySQL com o Hologres, consulte Use a Hologres catalog.
source:
type: mysql
using.built-in-catalog: mysql-catalog
tables: dbmysql.mysqltable
server-id: "8001-8004"
#(Optional) Synchronize data from newly added tables during the incremental phase.
scan.binlog.newly-added-table.enabled: true
#(Optional) Synchronize table and field comments.
include-comments.enabled: true
#(Optional) Prioritize the distribution of unbounded shards to prevent potential Out Of Memory (OOM) issues on TaskManagers.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(Optional) Enable parsing and filtering to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: hologres
using.built-in-catalog: hologres-catalog
jdbcWriteBatchSize: 1024 # Optional. Specify parameters for the sink table.
route:
- source-table: dbmysql.mysqltable
sink-table: public.holotable
Ler de uma tabela de dimensão do MySQL
INSERT INTO `<othersinktable>`
SELECT ...
FROM `<othersourcetable>` AS e
JOIN `<mysqlcatalog>`.`<dbname>`.`<tablename>` FOR SYSTEM_TIME AS OF e.proctime AS w
ON e.id = w.id;
Gravar em uma tabela do MySQL
INSERT INTO `<mysqlcatalog>`.`<dbname>`.`<tablename>`
SELECT ...
FROM `<othersourcetable>`