O conector do AnalyticDB for MySQL V3.0 permite ler e gravar dados em clusters AnalyticDB for MySQL V3.0 usando Flink SQL. O AnalyticDB for MySQL é um service de data warehouse nativo da cloud que suporta gravações em tempo real com alto throughput, análises de baixa latência e operações complexas de extração, transformação e carga (ETL).
Este conector oferece suporte aos seguintes tipos de tabela e capacidades:
| Item | Descrição |
|---|---|
| Tipos de tabela | Tabela source, tabela de dimensão e tabela sink. Nota
As tabelas source exigem o Ververica Runtime (VVR) 8.0.4 ou posterior. Para obter informações sobre os parâmetros da tabela source, consulte Use Flink to subscribe to binary logs. |
| Modos de execução | Modo streaming e modo batch |
| Formato de dados | N/A |
| Métricas | N/A |
| Tipo de API | API SQL |
| Atualização e exclusão de dados em tabelas sink | Suportado |
Pré-requisitos
Antes de começar, verifique se você possui:
Um cluster e uma tabela no AnalyticDB for MySQL. Consulte Create a cluster e CREATE TABLE
Uma lista de permissões de endereços IP configurada para o cluster. Consulte Configure an IP address whitelist
Sintaxe
CREATE TEMPORARY TABLE adb_table (
`id` INT,
`num` BIGINT,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'adb3.0',
'url' = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
'userName' = '<yourUsername>',
'password' = '<yourPassword>',
'tableName' = '<yourTablename>'
);
A chave primária definida na instrução DDL do Flink deve corresponder exatamente à chave primária da tabela física no AnalyticDB for MySQL — mesmos campos, mesmos nomes, presentes em ambas. Chaves primárias incompatíveis podem causar corrupção de dados.
Comportamentos principais
Chave primária e modo de gravação
Quando uma chave primária é definida na DDL, o conector opera no modo upsert, em que chaves primárias duplicadas acionam uma atualização em vez de uma nova inserção. O comportamento exato do SQL depende do parâmetro replaceMode:
replace— utilizaREPLACE INTO. Uma chave primária duplicada sobrescreve toda a linha existente.upsert— utilizaINSERT INTO ... ON DUPLICATE KEY UPDATE. Apenas os campos especificados são atualizados; os demais mantêm seus valores atuais.insert— utilizaINSERT IGNORE INTO. Chaves primárias duplicadas são ignoradas silenciosamente; a linha existente é preservada.
Se nenhuma chave primária for definida, o conector sempre usará INSERT IGNORE INTO.
O uso dereplaceModerequer o AnalyticDB for MySQL V3.1.3.5 ou posterior e o VVR 11.2 ou posterior para os valores de string (replace,upsert,insert). Versões anteriores do VVR utilizamtrue(equivalente areplace) efalse(equivalente aupsert). O VVR 11.2 e versões posteriores permanecem compatíveis comtrueefalse.
Buffer de gravação
O conector armazena registros em buffer na memória e os libera em lotes. A liberação ocorre quando uma das seguintes condições é atendida:
A quantidade de registros em buffer atinge
batchSize(padrão: 1.000) oubufferSize(padrão: 1.000)O tempo decorrido desde a última liberação atinge
flushIntervalMs(padrão: 3.000 ms)
Os parâmetros batchSize e bufferSize só têm efeito quando uma chave primária está definida.
Cache de tabela de dimensão
O conector suporta três políticas de cache para tabelas de dimensão. O uso de cache reduz as consultas à tabela física, mas exige memória adicional.
|
Política de cache |
Comportamento |
Indicado quando |
|
|
Todos os dados da tabela de dimensão são carregados na memória antes do início do job. As consultas subsequentes acessam apenas o cache. Os dados são recarregados após a expiração de |
A tabela de dimensão é pequena e a ausência de chaves é frequente |
|
|
As linhas acessadas com mais frequência são armazenadas em cache. Em caso de falha no cache, o conector consulta a tabela física e atualiza o cache. As entradas expiram após |
A tabela de dimensão é grande, mas o acesso se concentra em um subconjunto de chaves |
|
|
Sem cache. Cada consulta acessa diretamente a tabela física. |
Baixo volume de consultas ou requisitos rigorosos de atualização de dados |
Ao utilizar o cache ALL, o conector carrega de forma assíncrona toda a tabela de dimensão na memória. É necessário aumentar a memória do nó usado para junção de tabelas — o aumento deve ser de pelo menos o dobro do tamanho da tabela remota — para evitar erros de falta de memória (OOM). Monitore o uso de memória do nó durante a execução do job.
Parâmetros na cláusula WITH
Parâmetros comuns
|
Parâmetro |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
String |
Sim |
— |
Defina como |
|
|
String |
Sim |
— |
URL JDBC (Java Database Connectivity) do banco de dados, no formato |
|
|
String |
Sim |
— |
Nome de usuário para acesso ao banco de dados |
|
|
String |
Sim |
— |
Senha para acesso ao banco de dados |
|
|
String |
Sim |
— |
Nome da tabela de destino no banco de dados |
|
|
Integer |
Não |
10 |
Número máximo de tentativas em caso de falha de leitura ou gravação |
Parâmetros da tabela sink
|
Parâmetro |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
Integer |
Não |
1000 |
Quantidade de registros gravados por lote. Tem efeito apenas quando uma chave primária está definida. |
|
|
Integer |
Não |
1000 |
Quantidade máxima de registros armazenados em buffer antes de acionar a liberação. Tem efeito apenas quando uma chave primária está definida. |
|
|
Integer |
Não |
3000 |
Tempo máximo (em milissegundos) entre liberações. Ao atingir esse intervalo, todos os registros em buffer são gravados independentemente do tamanho do buffer. |
|
|
Boolean |
Não |
false |
Define se as operações de exclusão devem ser ignoradas. Defina como |
|
|
Boolean |
Não |
|
Modo de gravação quando uma chave primária está definida. VVR 11.2+: |
|
|
String |
Não |
(vazio) |
Lista separada por vírgulas de colunas a serem excluídas das atualizações quando |
|
|
Integer |
Não |
40 |
Número máximo de conexões simultâneas ao banco de dados no pool de threads |
Parâmetros da tabela de dimensão
|
Parâmetro |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
String |
Não |
ALL |
Política de cache: |
|
|
Integer |
Não |
100000 |
Quantidade máxima de linhas em cache. Obrigatório quando |
|
|
Integer |
Não |
Long.MAX_VALUE |
Tempo de vida das entradas de cache em milissegundos. Para |
|
|
Integer |
Não |
1024 |
Quantidade máxima de linhas da tabela de dimensão correspondentes por registro de entrada. Se cada registro de entrada corresponder a no máximo n linhas da tabela de dimensão, defina este valor como n para otimizar o desempenho da junção no Realtime Compute for Apache Flink. |
Mapeamento de tipos de dados
|
AnalyticDB for MySQL V3.0 |
Realtime Compute for Apache Flink |
|
BOOLEAN |
BOOLEAN |
|
TINYINT |
TINYINT |
|
SMALLINT |
SMALLINT |
|
INT |
INT |
|
BIGINT |
BIGINT |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
DECIMAL(p, s) ou NUMERIC(p, s) |
DECIMAL(p, s) |
|
VARCHAR |
STRING |
|
BINARY |
BYTES |
|
DATE |
DATE |
|
TIME |
TIME |
|
DATETIME |
TIMESTAMP |
|
TIMESTAMP |
TIMESTAMP |
|
POINT |
STRING |
Exemplos
Tabela sink
O exemplo a seguir lê dados de uma source datagen e grava os dados em uma tabela sink do AnalyticDB for MySQL.
CREATE TEMPORARY TABLE datagen_source (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE adb_sink (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'adb3.0',
'url' = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
'userName' = '<yourUsername>',
'password' = '<yourPassword>',
'tableName' = '<yourTablename>'
);
INSERT INTO adb_sink
SELECT * FROM datagen_source;
Tabela de dimensão
O exemplo abaixo realiza uma junção temporal entre uma source datagen e uma tabela de dimensão do AnalyticDB for MySQL. Os resultados são gravados em uma sink blackhole.
CREATE TEMPORARY TABLE datagen_source (
`a` INT,
`b` VARCHAR,
`c` STRING,
`proctime` AS PROCTIME()
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE adb_dim (
`a` INT,
`b` VARCHAR,
`c` VARCHAR
) WITH (
'connector' = 'adb3.0',
'url' = 'jdbc:mysql://<endpoint>:<port>/<databaseName>',
'userName' = '<yourUsername>',
'password' = '<yourPassword>',
'tableName' = '<yourTablename>'
);
CREATE TEMPORARY TABLE blackhole_sink (
`a` INT,
`b` VARCHAR
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN adb_dim FOR SYSTEM_TIME AS OF T.proctime AS H
ON T.a = H.a;