O conector JDBC lê e escreve em bancos de dados relacionais — MySQL, PostgreSQL e Oracle — usando DDL SQL padrão em jobs do Flink SQL.
Tipos de tabela suportados: tabela source · tabela de dimensão · tabela sink
Modos de execução suportados: modo streaming · modo batch · API SQL
Pré-requisitos
Antes de começar, verifique se:
O banco de dados e a tabela de destino já existem
O JAR do driver JDBC do seu banco de dados está disponível para upload
Limitações
Leituras limitadas: Uma tabela source JDBC é uma source limitada. A tarefa termina após a leitura de todas as linhas. Para capturar dados de alteração em tempo real, use um conector Change Data Capture (CDC) — consulte Crie uma tabela source CDC do MySQL e Crie uma tabela source CDC do PostgreSQL (visualização pública).
Versão do PostgreSQL: A escrita no PostgreSQL exige a versão 9.5 ou posterior, pois a sink depende da cláusula
ON CONFLICT.-
Upload do driver JDBC: Faça o upload do JAR do driver JDBC como arquivo de dependência antes de executar o job. Drivers comuns:
Banco de dados
Group ID
Artifact ID
MySQL
mysqlOracle
com.oracle.database.jdbcPostgreSQL
org.postgresqlPara drivers não listados aqui, verifique a compatibilidade antes do uso.
-
Comportamento de upsert do MySQL: Ao escrever em uma tabela sink do MySQL com chave primária, o conector emite instruções
INSERT INTO ... ON DUPLICATE KEY UPDATE ....AvisoA inserção de linhas com valores de índice exclusivo duplicados — mesmo quando as chaves primárias diferem — sobrescreve as linhas existentes em qualquer tabela física com restrição de índice exclusivo, causando perda de dados.
Crie uma tabela JDBC
CREATE TABLE jdbc_table (
`id` BIGINT,
`name` VARCHAR,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:<db-type>://<host>:<port>/<database>',
'table-name' = '<yourTable>',
'username' = '<yourUsername>',
'password' = '<yourPassword>'
);
Opções do conector
Geral
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
STRING |
Sim |
— |
Defina como |
|
|
STRING |
Sim |
— |
URL JDBC do banco de dados. |
|
|
STRING |
Sim |
— |
Nome da tabela de leitura ou escrita. |
|
|
STRING |
Não |
— |
Nome de usuário do banco de dados. Configure junto com |
|
|
STRING |
Não |
— |
Senha do banco de dados. |
Opções de source
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
STRING |
Não |
— |
Coluna usada para dividir os dados em partições. Deve ser do tipo NUMERIC ou TIMESTAMP. Consulte Varredura particionada. |
|
|
INTEGER |
Não |
— |
Número de partições. |
|
|
LONG |
Não |
— |
Menor valor da primeira partição. |
|
|
LONG |
Não |
— |
Maior valor da última partição. |
|
|
INTEGER |
Não |
|
Linhas buscadas por ida e volta ao banco de dados. Se definido como |
|
|
BOOLEAN |
Não |
|
Ative o auto-commit para transações de leitura. |
Opções de sink
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
INTEGER |
Não |
|
Máximo de linhas no buffer antes do flush. Defina como |
|
|
DURATION |
Não |
|
Tempo máximo de retenção das linhas no buffer antes do flush. Defina como |
|
|
INTEGER |
Não |
|
Número máximo de tentativas de escrita em caso de falha. |
|
|
BOOLEAN |
Não |
|
Ignora mensagens de exclusão em vez de encaminhá-las. Requer VVR 11.4+. |
|
|
STRING |
Não |
|
Controla quais mensagens de exclusão são ignoradas quando |
Para liberar assincronamente as linhas do buffer por temporizador, em vez de por contagem de linhas, defina sink.buffer-flush.max-rows como 0 e configure sink.buffer-flush.interval com o intervalo de flush desejado.
Opções de tabela de dimensão
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
INTEGER |
Não |
— |
Número máximo de linhas no cache de lookup. Quando o cache enche, a linha menos recentemente usada expira. O cache permanece desativado a menos que tanto |
|
|
DURATION |
Não |
— |
Tempo máximo de validade de uma linha em cache antes da expiração. |
|
|
BOOLEAN |
Não |
|
Armazena resultados vazios de lookup em cache para evitar consultas repetidas ao banco de dados para chaves inexistentes. |
|
|
INTEGER |
Não |
|
Número máximo de tentativas quando uma consulta ao banco de dados falha. |
Opções do PostgreSQL
|
Opção |
Tipo |
Obrigatório |
Padrão |
Descrição |
|
|
BOOLEAN |
Não |
|
Quando definido como |
Comportamentos principais
Varredura particionada
Para ativar leituras paralelas de uma tabela source, configure as opções scan.partition.column, scan.partition.lower-bound e scan.partition.upper-bound junto com scan.partition.num. A coluna de partição deve ser numérica ou TIMESTAMP. Para detalhes sobre o cálculo das divisões, consulte Partitioned Scan na documentação do Apache Flink.
Cache de lookup
Por padrão, toda consulta de join de lookup acessa o banco de dados diretamente. Ative o cache LRU definindo tanto lookup.cache.max-rows quanto lookup.cache.ttl para reduzir a carga no banco de dados, aceitando dados ligeiramente desatualizados como contrapartida.
Escritas idempotentes
Quando a tabela sink possui chave primária, o conector emite instruções de upsert específicas do banco de dados. Para MySQL, o conector usa INSERT ... ON DUPLICATE KEY UPDATE ....
Exemplos
Os três exemplos usam um conector blackhole ou datagen como contraparte leve para que as instruções sejam autossuficientes.
Leitura de um banco de dados (tabela source)
CREATE TEMPORARY TABLE jdbc_source (
`id` INT,
`name` VARCHAR
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/mydb',
'table-name' = '<yourTable>',
'username' = '<yourUsername>',
'password' = '<yourPassword>'
);
CREATE TEMPORARY TABLE blackhole_sink (
`id` INT,
`name` VARCHAR
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink SELECT * FROM jdbc_source;
Escrita em um banco de dados (tabela sink)
CREATE TEMPORARY TABLE datagen_source (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE jdbc_sink (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/mydb',
'table-name' = '<yourTable>',
'username' = '<yourUsername>',
'password' = '<yourPassword>'
);
INSERT INTO jdbc_sink SELECT * FROM datagen_source;
Enriquecimento de stream com lookups de banco de dados (tabela de dimensão)
CREATE TEMPORARY TABLE datagen_source (
`id` INT,
`data` BIGINT,
`proctime` AS PROCTIME()
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE jdbc_dim (
`id` INT,
`name` VARCHAR
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/mydb',
'table-name' = '<yourTable>',
'username' = '<yourUsername>',
'password' = '<yourPassword>'
);
CREATE TEMPORARY TABLE blackhole_sink (
`id` INT,
`data` BIGINT,
`name` VARCHAR
) WITH (
'connector' = 'blackhole'
);
-- Look up the matching name for each stream record at processing time
INSERT INTO blackhole_sink
SELECT T.`id`, T.`data`, H.`name`
FROM datagen_source AS T
JOIN jdbc_dim FOR SYSTEM_TIME AS OF T.proctime AS H
ON T.id = H.id;
Mapeamentos de tipos de dados
|
MySQL |
Oracle |
PostgreSQL |
Flink SQL |
|
|
— |
— |
|
|
|
— |
|
|
|
|
— |
|
|
|
|
— |
|
|
|
|
— |
— |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
— |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
— |
— |
|
|
Próximos passos
Crie uma tabela source CDC do MySQL — capture dados de alteração em tempo real do MySQL
Crie uma tabela source CDC do PostgreSQL (visualização pública) — capture dados de alteração em tempo real do PostgreSQL