Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Connecteur JDBC

Dernière mise à jour :Aug 09, 2026

Le connecteur JDBC lit et écrit dans des bases de données relationnelles (MySQL, PostgreSQL et Oracle) en utilisant le langage DDL SQL standard au sein des jobs Flink SQL.

Types de table pris en charge : table source · table de dimension · table de sortie (sink)

Modes d'exécution pris en charge : mode streaming · mode batch · API SQL

Prérequis

Avant de commencer, assurez-vous que :

  • La base de données et la table cibles existent déjà.

  • Le fichier JAR du pilote JDBC correspondant à votre base de données est disponible pour téléchargement.

Limites

  • Lectures bornées : Une table source JDBC constitue une source bornée. La tâche se termine une fois toutes les lignes lues. Pour capturer les données de modification en temps réel, utilisez plutôt un connecteur Change Data Capture (CDC). Consultez les rubriques Créer une table source MySQL CDC et Créer une table source PostgreSQL CDC (aperçu public).

  • Version de PostgreSQL : L'écriture dans PostgreSQL nécessite la version 9.5 ou ultérieure, car le mécanisme de sortie (sink) s'appuie sur la clause ON CONFLICT.

  • Téléchargement du pilote JDBC : Téléchargez le fichier JAR du pilote JDBC en tant que dépendance avant d'exécuter votre job. Pilotes courants :

    Base de données Group ID Artifact ID
    MySQL mysql mysql-connector-java
    Oracle com.oracle.database.jdbc ojdbc8
    PostgreSQL org.postgresql postgresql

    Pour les pilotes non répertoriés ici, vérifiez la compatibilité avant utilisation.

  • Comportement upsert MySQL : Lors de l'écriture dans une table de sortie MySQL dotée d'une clé primaire, le connecteur émet des instructions INSERT INTO ... ON DUPLICATE KEY UPDATE ....

    Avertissement

    L'insertion de lignes comportant des valeurs d'index unique dupliquées, même si les clés primaires diffèrent, entraîne l'écrasement des lignes existantes dans toute table physique soumise à une contrainte d'index unique, ce qui provoque une perte de données.

Créer une table 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>'
);

Options du connecteur

Général

Option Type Obligatoire Par défaut Description
connector STRING Oui Définissez cette option sur jdbc.
url STRING Oui URL JDBC de la base de données.
table-name STRING Oui Nom de la table à lire ou à écrire.
username STRING Non Nom d'utilisateur de la base de données. À définir conjointement avec password.
password STRING Non Mot de passe de la base de données.

Options de source

Option Type Obligatoire Par défaut Description
scan.partition.column STRING Non Colonne utilisée pour diviser les données en partitions. Doit être de type NUMERIC ou TIMESTAMP. Consultez la section Analyse partitionnée.
scan.partition.num INTEGER Non Nombre de partitions.
scan.partition.lower-bound LONG Non Valeur minimale de la première partition.
scan.partition.upper-bound LONG Non Valeur maximale de la dernière partition.
scan.fetch-size INTEGER Non 0 Nombre de lignes récupérées par aller-retour vers la base de données. Si cette option est définie sur 0, elle est ignorée.
scan.auto-commit BOOLEAN Non true Active commit automatique pour les transactions de lecture.

Options de sortie (sink)

Option Type Obligatoire Par défaut Description
sink.buffer-flush.max-rows INTEGER Non 100 Nombre maximal de lignes mises en mémoire tampon avant vidage. Définissez cette option sur 0 pour vider chaque ligne immédiatement.
sink.buffer-flush.interval DURATION Non 1000 ms Durée maximale de conservation des lignes en mémoire tampon avant vidage. Définissez cette option sur 0 pour vider chaque ligne immédiatement.
sink.max-retries INTEGER Non 3 Nombre maximal de tentatives d'écriture en cas d'échec.
sink.ignore-delete BOOLEAN Non false Ignore les messages de suppression au lieu de les transmettre. Nécessite VVR 11.4 ou version ultérieure.
sink.ignore-delete-mode STRING Non ALL Contrôle les messages de suppression ignorés lorsque sink.ignore-delete est défini sur true. Nécessite VVR 11.4 ou version ultérieure. Valeurs valides : ALL (ignore -D et -U), REAL_DELETE (ignore uniquement -D), UPDATE_BEFORE (ignore uniquement -U).
Remarque

Pour vider les lignes mises en mémoire tampon de manière asynchrone selon un minuteur plutôt que sur la base du nombre de lignes, définissez sink.buffer-flush.max-rows sur 0 et configurez sink.buffer-flush.interval selon l'intervalle de vidage souhaité.

Options de table de dimension

Option Type Obligatoire Par défaut Description
lookup.cache.max-rows INTEGER Non Nombre maximal de lignes dans le cache de recherche. Lorsque le cache est plein, la ligne la moins récemment utilisée expire. La mise en cache est désactivée sauf si les options lookup.cache.max-rows et lookup.cache.ttl sont toutes deux définies.
lookup.cache.ttl DURATION Non Durée maximale de validité d'une ligne mise en cache avant son expiration.
lookup.cache.caching-missing-key BOOLEAN Non true Met en cache les résultats de recherche vides afin d'éviter des interrogations répétées de la base de données pour les clés manquantes.
lookup.max-retries INTEGER Non 3 Nombre maximal de tentatives en cas d'échec d'une requête vers la base de données.

Options PostgreSQL

Option Type Obligatoire Par défaut Description
source.extend-type.enabled BOOLEAN Non false Lorsque cette option est définie sur true, elle mappe les colonnes PostgreSQL JSONB et UUID vers le type Flink STRING. Pour les jointures de recherche sur des colonnes UUID, ajoutez également stringtype=unspecified à l'URL JDBC afin que PostgreSQL effectue les requêtes selon le type réel plutôt que par conversion.

Comportements clés

Analyse partitionnée

Pour activer les lectures parallèles depuis une table source, configurez conjointement les options scan.partition.column, scan.partition.lower-bound, scan.partition.upper-bound et scan.partition.num. La colonne de partition doit être de type numérique ou TIMESTAMP. Pour plus de détails sur le calcul des splits, consultez la documentation Apache Flink relative à l'analyse partitionnée.

Cache de recherche

Par défaut, chaque requête de jointure de recherche interroge directement la base de données. Activez le cache LRU en définissant simultanément lookup.cache.max-rows et lookup.cache.ttl afin de réduire la charge sur la base de données, au prix d'une légère obsolescence des données.

Écritures idempotentes

Lorsque la table de sortie possède une clé primaire, le connecteur émet des instructions upsert spécifiques à la base de données. Pour MySQL, le connecteur utilise l'instruction INSERT ... ON DUPLICATE KEY UPDATE ....

Exemples

Les trois exemples suivants utilisent un connecteur blackhole ou datagen comme contrepartie légère, permettant ainsi aux instructions d'être autonomes.

Lire depuis une base de données (table 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;

Écrire dans une base de données (table de sortie)

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;

Enrichir un flux via des recherches en base de données (table de dimension)

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;

Mappings de types de données

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

Étapes suivantes