Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector JDBC

Última atualização: Jun 27, 2026

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

    mysql

    mysql-connector-java

    Oracle

    com.oracle.database.jdbc

    ojdbc8

    PostgreSQL

    org.postgresql

    postgresql

    Para 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 ....

    Aviso

    A 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

connector

STRING

Sim

Defina como jdbc.

url

STRING

Sim

URL JDBC do banco de dados.

table-name

STRING

Sim

Nome da tabela de leitura ou escrita.

username

STRING

Não

Nome de usuário do banco de dados. Configure junto com password.

password

STRING

Não

Senha do banco de dados.

Opções de source

Opção

Tipo

Obrigatório

Padrão

Descrição

scan.partition.column

STRING

Não

Coluna usada para dividir os dados em partições. Deve ser do tipo NUMERIC ou TIMESTAMP. Consulte Varredura particionada.

scan.partition.num

INTEGER

Não

Número de partições.

scan.partition.lower-bound

LONG

Não

Menor valor da primeira partição.

scan.partition.upper-bound

LONG

Não

Maior valor da última partição.

scan.fetch-size

INTEGER

Não

0

Linhas buscadas por ida e volta ao banco de dados. Se definido como 0, esta opção é ignorada.

scan.auto-commit

BOOLEAN

Não

true

Ative o auto-commit para transações de leitura.

Opções de sink

Opção

Tipo

Obrigatório

Padrão

Descrição

sink.buffer-flush.max-rows

INTEGER

Não

100

Máximo de linhas no buffer antes do flush. Defina como 0 para liberar cada linha imediatamente.

sink.buffer-flush.interval

DURATION

Não

1000 ms

Tempo máximo de retenção das linhas no buffer antes do flush. Defina como 0 para liberar cada linha imediatamente.

sink.max-retries

INTEGER

Não

3

Número máximo de tentativas de escrita em caso de falha.

sink.ignore-delete

BOOLEAN

Não

false

Ignora mensagens de exclusão em vez de encaminhá-las. Requer VVR 11.4+.

sink.ignore-delete-mode

STRING

Não

ALL

Controla quais mensagens de exclusão são ignoradas quando sink.ignore-delete é true. Requer VVR 11.4+. Valores válidos: ALL (ignora -D e -U), REAL_DELETE (ignora apenas -D), UPDATE_BEFORE (ignora apenas -U).

Nota

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

lookup.cache.max-rows

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 lookup.cache.max-rows quanto lookup.cache.ttl estejam definidos.

lookup.cache.ttl

DURATION

Não

Tempo máximo de validade de uma linha em cache antes da expiração.

lookup.cache.caching-missing-key

BOOLEAN

Não

true

Armazena resultados vazios de lookup em cache para evitar consultas repetidas ao banco de dados para chaves inexistentes.

lookup.max-retries

INTEGER

Não

3

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

source.extend-type.enabled

BOOLEAN

Não

false

Quando definido como true, mapeia colunas JSONB e UUID do PostgreSQL para STRING do Flink. Para joins de lookup em colunas UUID, adicione também stringtype=unspecified à URL JDBC para que o PostgreSQL consulte pelo tipo real em vez de converter o tipo (casting).

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

TINYINT

TINYINT

SMALLINT, TINYINT UNSIGNED

SMALLINT, INT2, SMALLSERIAL, SERIAL2

SMALLINT

INT, MEDIUMINT, SMALLINT UNSIGNED

INTEGER, SERIAL

INT

BIGINT, INT UNSIGNED

BIGINT, BIGSERIAL

BIGINT

BIGINT UNSIGNED

DECIMAL(20, 0)

FLOAT

BINARY_FLOAT

REAL, FLOAT4

FLOAT

DOUBLE, DOUBLE PRECISION

BINARY_DOUBLE

FLOAT8, DOUBLE PRECISION

DOUBLE

NUMERIC(p, s), DECIMAL(p, s)

SMALLINT, FLOAT(s), DOUBLE PRECISION, REAL, NUMBER(p, s)

NUMERIC(p, s), DECIMAL(p, s)

DECIMAL(p, s)

BOOLEAN, TINYINT(1)

BOOLEAN

BOOLEAN

DATE

DATE

DATE

DATE

TIME [(p)]

DATE

TIME [(p)] [WITHOUT TIMEZONE]

TIME [(p)] [WITHOUT TIMEZONE]

DATETIME [(p)]

TIMESTAMP [(p)] [WITHOUT TIMEZONE]

TIMESTAMP [(p)] [WITHOUT TIMEZONE]

TIMESTAMP [(p)] [WITHOUT TIMEZONE]

CHAR(n), VARCHAR(n), TEXT

CHAR(n), VARCHAR(n), CLOB

CHAR(n), CHARACTER(n), VARCHAR(n), CHARACTER VARYING(n), TEXT, JSONB, UUID

STRING

BINARY, VARBINARY, BLOB

RAW(s), BLOB

BYTEA

BYTES

ARRAY

ARRAY

Próximos passos