Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:ApsaraDB RDS for MySQL

Última atualização: Aug 20, 2026
Importante

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:

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

NONE

Sem cache — toda consulta acessa diretamente o banco de dados

Requisitos de baixa latência e conjuntos de dados pequenos

LRU

Armazena em cache um número fixo de linhas usadas recentemente por task manager

Subconjuntos de tabelas grandes acessados frequentemente

ALL

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'
);
Nota

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

connector

STRING

Sim

Defina como rds

tableName

STRING

Sim

Nome da tabela física no ApsaraDB RDS for MySQL

userName

STRING

Sim

Nome de usuário do banco de dados

password

STRING

Sim

Senha do banco de dados

url

STRING

Sim

Endpoint da Virtual Private Cloud (VPC) do banco de dados, no formato jdbc:mysql://<internal-endpoint>:<port>/<database-name>. Para tabelas de destino, adicione ?rewriteBatchedStatements=true. Para detalhes sobre endpoints, consulte Visualize and change the internal and public endpoints and port numbers of an ApsaraDB RDS for MySQL instance

maxRetryTimes

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

batchSize

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

bufferSize

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

flushIntervalMs

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 batchSize ou bufferSize terem sido atingidos. Se definido como 0 (padrão para VVR 4.0.0–4.0.6), pequenos volumes de dados em buffer podem nunca ser gravados — atualize para uma versão mais recente do VVR para evitar isso

ignoreDelete

BOOLEAN

Não

false

Defina como true para ignorar operações de exclusão. Útil quando múltiplos operadores atualizam campos diferentes da mesma linha — sem essa configuração, uma exclusão em um operador seguida de uma atualização parcial em outro deixa os campos não atualizados como null ou com seus valores padrão

connectionMaxActive

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

cache

STRING

Não

NONE (VVR anterior a 4.0.6), ALL (VVR 4.0.6+)

Política de cache. Valores válidos: NONE, LRU, ALL. Consulte Cache policies

cacheSize

INTEGER

Não

100000

Número máximo de linhas a serem armazenadas em cache. Obrigatório quando cache estiver definido como LRU; ignorado para NONE e ALL

cacheTTLMs

LONG

Não

Sem expiração para NONE e LRU; sem recarregamento para ALL

Tempo de vida do cache em milissegundos. Para LRU, as linhas expiram após esse período. Para ALL, todo o cache é recarregado neste intervalo

maxJoinRows

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

numRecordsOut

Total de linhas gravadas

numRecordsOutPerSecond

Linhas gravadas por segundo

numBytesOut

Total de bytes gravados

numBytesOutPerSecond

Bytes gravados por segundo

currentSendTime

Latência atual de gravação

numRecordsOutErrors

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