O conector do ApsaraDB RDS for MySQL não terá suporte futuro. Use o MySQL connector como alternativa.
O conector do ApsaraDB RDS for MySQL permite gravar a saída do Flink SQL em uma tabela de destino do ApsaraDB RDS for MySQL ou associar um stream a uma tabela de dimensão do ApsaraDB RDS for MySQL.
Tipos de tabela suportados: Tabela de destino · Tabela de dimensão
Supported running modes: Modo em lote · Modo streaming
Tipo de API: SQL
Atualizações e exclusões de dados em tabelas de destino: Suportadas
Pré-requisitos
Antes de começar, verifique se você tem:
Um banco de dados e uma tabela no ApsaraDB RDS for MySQL. Consulte Crie databases and accounts for an ApsaraDB RDS for MySQL instance
Uma lista de permissões de endereços IP configurada para o banco de dados. Consulte Connect to an ApsaraDB RDS for MySQL instance using a database client or the CLI
Limitações
Requer o Realtime Compute for Apache Flink com Ververica Runtime (VVR) 2.0.0 ou superior. Para melhor desempenho e estabilidade, use o VVR 6.X ou posterior.
Apenas bancos de dados ApsaraDB RDS for MySQL são suportados.
O conector usa semântica at-least-once. Se a tabela de destino tiver chave primária, a idempotência garante a integridade dos dados.
Funcionamento
Comportamento de gravação na tabela de destino
Cada linha de saída é convertida em uma instrução SQL antes da gravação na tabela de destino:
Sem chave primária — executa
INSERT INTO table_name (col1, col2, ...) VALUES (val1, val2, ...);Com chave primária — executa
INSERT INTO table_name (col1, col2, ...) VALUES (val1, val2, ...) ON DUPLICATE KEY UPDATE col1 = VALUES(col1), col2 = VALUES(col2), ...;
Conflitos de índice único: Se a tabela física tiver uma restrição de índice único além da chave primária, inserir duas linhas com chaves primárias diferentes, mas com o mesmo valor de índice único, sobrescreve a linha anterior e causa perda de dados.
Chaves primárias com incremento automático: Não declare campos de incremento automático na DDL do Flink. O banco de dados atribui esses valores automaticamente. O conector pode gravar e excluir linhas com campos de incremento automático, mas não pode atualizá-las.
Políticas de cache para tabelas de dimensão
O conector oferece três políticas de cache para consultas em tabelas de dimensão:
|
Política |
Comportamento |
Quando usar |
|
|
Sem cache — toda consulta acessa diretamente o banco de dados |
Requisitos de baixa latência e conjuntos de dados pequenos |
|
|
Armazena em cache um número fixo de linhas usadas recentemente por task manager |
Subconjuntos de tabelas grandes acessados frequentemente |
|
|
Carrega toda a tabela na memória e a recarrega periodicamente |
Tabelas de referência pequenas e estáticas |
Sintaxe
Tabela de destino
CREATE TABLE rds_sink (
id INT,
num BIGINT,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'rds',
'tableName' = '<your-table-name>',
'userName' = '<your-user-name>',
'password' = '<your-password>',
'url' = 'jdbc:mysql://<internal-endpoint>:<port>/<database-name>?rewriteBatchedStatements=true'
);
Adicione ?rewriteBatchedStatements=true ao valor de url nas tabelas de destino para aumentar o throughput de gravação.
Tabela de dimensão
CREATE TABLE rds_dim (
id1 INT,
id2 VARCHAR
) WITH (
'connector' = 'rds',
'tableName' = '<your-table-name>',
'userName' = '<your-user-name>',
'password' = '<your-password>',
'url' = 'jdbc:mysql://<internal-endpoint>:<port>/<database-name>',
'cache' = 'NONE'
);
Parâmetros na cláusula WITH
Parâmetros comuns
|
Parâmetro |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
STRING |
Sim |
— |
Defina como |
|
|
STRING |
Sim |
— |
Nome da tabela física no ApsaraDB RDS for MySQL |
|
|
STRING |
Sim |
— |
Nome de usuário do banco de dados |
|
|
STRING |
Sim |
— |
Senha do banco de dados |
|
|
STRING |
Sim |
— |
Endpoint da Virtual Private Cloud (VPC) do banco de dados, no formato |
|
|
INTEGER |
Não |
10 (VVR 4.0.7+), 3 (VVR 4.0.6 e anteriores) |
Número máximo de tentativas para consultas falhas em tabelas de dimensão ou gravações na tabela de destino |
Parâmetros da tabela de destino
|
Parâmetro |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
INTEGER |
Não |
4096 (VVR 4.0.7+), 5000 (VVR 4.0.0–4.0.6), 100 (VVR 3.x e anteriores) |
Quantidade de linhas gravadas por lote |
|
|
INTEGER |
Não |
10000 |
Máximo de linhas armazenadas em cache na memória antes de acionar uma gravação. Suportado no VVR 4.0.7 e posteriores. Só tem efeito quando há uma chave primária definida |
|
|
INTEGER |
Não |
2000 (VVR 4.0.7+), 0 (VVR 4.0.0–4.0.6), 1000 (VVR 3.x e anteriores) |
Intervalo em milissegundos para liberar o buffer na tabela de destino, independentemente de os limiares de |
|
|
BOOLEAN |
Não |
false |
Defina como |
|
|
INTEGER |
Não |
40 |
Tamanho do pool de conexões. Suportado no VVR 4.0.7 e posteriores. Aumente este valor se ocorrerem timeouts no pool de conexões; diminua-o se o banco de dados limitar o número de conexões simultâneas |
Parâmetros da tabela de dimensão
|
Parâmetro |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
STRING |
Não |
NONE (VVR anterior a 4.0.6), ALL (VVR 4.0.6+) |
Política de cache. Valores válidos: |
|
|
INTEGER |
Não |
100000 |
Número máximo de linhas a serem armazenadas em cache. Obrigatório quando |
|
|
LONG |
Não |
Sem expiração para |
Tempo de vida do cache em milissegundos. Para |
|
|
INTEGER |
Não |
1024 |
Número máximo de linhas da tabela de dimensão correspondentes por linha de entrada. Configure este valor com o máximo esperado de linhas de dimensão por linha da tabela principal para evitar varreduras desnecessárias |
Métricas
A tabela de destino expõe as seguintes métricas. Tabelas de dimensão não possuem métricas.
|
Métrica |
Descrição |
|
|
Total de linhas gravadas |
|
|
Linhas gravadas por segundo |
|
|
Total de bytes gravados |
|
|
Bytes gravados por segundo |
|
|
Latência atual de gravação |
|
|
Total de erros de gravação |
Para definições das métricas, consulte Métricas.
Mapeamentos de tipos de dados
|
Tipo Flink |
Tipo ApsaraDB RDS for MySQL |
|
BOOLEAN |
BOOLEAN |
|
TINYINT |
TINYINT |
|
TINYINT(1) (apenas tabelas de dimensão) |
BOOLEAN |
|
SMALLINT |
SMALLINT |
|
SMALLINT |
TINYINT UNSIGNED |
|
INT |
INT |
|
INT |
SMALLINT UNSIGNED |
|
BIGINT |
BIGINT |
|
BIGINT |
INT UNSIGNED |
|
DECIMAL(20, 0) |
BIGINT UNSIGNED |
|
FLOAT |
FLOAT |
|
DECIMAL |
DECIMAL |
|
DOUBLE |
DOUBLE |
|
DATE |
DATE |
|
TIME |
TIME |
|
TIMESTAMP |
TIMESTAMP |
|
VARCHAR |
VARCHAR |
|
VARBINARY |
VARBINARY |
Exemplos
Exemplo de tabela de destino
O exemplo abaixo lê dados de uma source DataGen e grava em uma tabela de destino do ApsaraDB RDS for MySQL.
CREATE TEMPORARY TABLE datagen_source (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE rds_sink (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'rds',
'tableName' = '<your-table-name>',
'userName' = '<your-user-name>',
'password' = '<your-password>',
'url' = 'jdbc:mysql://<internal-endpoint>:<port>/<database-name>?rewriteBatchedStatements=true'
);
INSERT INTO rds_sink
SELECT * FROM datagen_source;
Exemplo de tabela de dimensão
Este exemplo associa um stream a uma tabela de dimensão do ApsaraDB RDS for MySQL usando um temporal join.
CREATE TEMPORARY TABLE datagen_source (
a INT,
b BIGINT,
c STRING,
`proctime` AS PROCTIME()
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE rds_dim (
a INT,
b VARCHAR,
c VARCHAR
) WITH (
'connector' = 'rds',
'tableName' = '<your-table-name>',
'userName' = '<your-user-name>',
'password' = '<your-password>',
'url' = 'jdbc:mysql://<internal-endpoint>:<port>/<database-name>'
);
CREATE TEMPORARY TABLE blackhole_sink (
a INT,
b STRING
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN rds_dim FOR SYSTEM_TIME AS OF T.`proctime` AS H
ON T.a = H.a;
Perguntas frequentes
Próximos passos
MySQL connector — substituto recomendado para este conector
ApsaraDB RDS for MySQL — visão geral do product e documentação de recursos
Métricas — definições para todas as métricas do conector